mirror of
https://github.com/Portabase/agent.git
synced 2026-10-03 21:03:16 +00:00
Compare commits
51 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| d42c31f3dc | |||
| 84673373a7 | |||
| 6f8a5c4228 | |||
| c5806a6ea4 | |||
| c241020542 | |||
| e0bc6556aa | |||
| 947cb9bed7 | |||
| b5d0035a3b | |||
| d52c328a99 | |||
| 9b79ac4850 | |||
| 6e50718b65 | |||
| ee5ed81d48 | |||
| df078b6202 | |||
| 24e3e3098d | |||
| f6e1a0df98 | |||
| 49366ee24e | |||
| 012a984973 | |||
| 6d1d7c6891 | |||
| c70cde8a9c | |||
| c353443bdd | |||
| 9f9dedfe87 | |||
| 9bb830327e | |||
| 31e55a7b4f | |||
| 576c055092 | |||
| ba83b6b2b7 | |||
| ea356ed54f | |||
| b222b13934 | |||
| 9409f0ff25 | |||
| 59312215e9 | |||
| 3ff360a971 | |||
| cc1703592e | |||
| 5f92297108 | |||
| d1821f7045 | |||
| 0fe4d50cbd | |||
| 1f29466a28 | |||
| b74aaa0bbc | |||
| 7b5e5b2c78 | |||
| 99f1ef2081 | |||
| 42c3c5945a | |||
| db0f87c2e9 | |||
| f95f0aa73a | |||
| 2177dfd44b | |||
| 0a53eec184 | |||
| 87f33af772 | |||
| b6a120fcbf | |||
| 5298d82576 | |||
| b0da2e40a3 | |||
| 55e20d48e7 | |||
| f7f5f7e141 | |||
| 2170f96a72 | |||
| 298d46ba81 |
+3
-5
@@ -1,11 +1,9 @@
|
||||
# Git
|
||||
.git
|
||||
.gitignore
|
||||
|
||||
# MD files
|
||||
CHANGELOG.md
|
||||
README.md
|
||||
RELEASE.md
|
||||
|
||||
#IDE configurations
|
||||
.idea
|
||||
target
|
||||
dump.rdb
|
||||
.superpowers
|
||||
|
||||
@@ -44,21 +44,6 @@ jobs:
|
||||
- name: Cache cargo build
|
||||
uses: Swatinem/rust-cache@v2
|
||||
|
||||
- name: Cache vcpkg installed packages
|
||||
uses: actions/cache@v4
|
||||
with:
|
||||
path: C:\vcpkg\installed
|
||||
key: vcpkg-openssl-x64-windows-v1
|
||||
|
||||
- name: Install OpenSSL (x64) via vcpkg
|
||||
shell: pwsh
|
||||
run: |
|
||||
# windows-latest ships vcpkg preinstalled; the install is a no-op when the
|
||||
# package is restored from cache.
|
||||
& "$env:VCPKG_INSTALLATION_ROOT\vcpkg.exe" install openssl:x64-windows
|
||||
'VCPKG_ROOT=C:\vcpkg' | Out-File -FilePath $env:GITHUB_ENV -Encoding utf8 -Append
|
||||
'OPENSSL_DIR=C:\vcpkg\installed\x64-windows' | Out-File -FilePath $env:GITHUB_ENV -Encoding utf8 -Append
|
||||
|
||||
- name: Build (cargo release)
|
||||
shell: pwsh
|
||||
run: cargo build --release --bin app
|
||||
|
||||
+2
-1
@@ -6,5 +6,6 @@
|
||||
.env
|
||||
|
||||
.claude
|
||||
|
||||
/docs
|
||||
|
||||
.superpowers
|
||||
|
||||
+1
-1
@@ -27,5 +27,5 @@ keywords:
|
||||
- self-hosted
|
||||
- portabase
|
||||
license: Apache-2.0
|
||||
version: 1.19.0
|
||||
version: 1.21.1
|
||||
date-released: '2026-02-24'
|
||||
|
||||
Generated
+1215
-1354
File diff suppressed because it is too large
Load Diff
+4
-3
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "portabase-agent"
|
||||
version = "1.19.0"
|
||||
version = "1.21.1"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
@@ -19,10 +19,11 @@ log = "0.4.29"
|
||||
toml = "0.9.10"
|
||||
reqwest = { version = "0.13.1", features = ["json", "blocking", "multipart", "stream", "query"] }
|
||||
anyhow = "1.0.100"
|
||||
tokio = { version = "1.49.0", features = ["rt", "rt-multi-thread", "macros", "fs"] }
|
||||
tokio = { version = "1.49.0", features = ["rt", "rt-multi-thread", "macros", "fs", "process", "io-util"] }
|
||||
async-trait = "0.1.89"
|
||||
tempfile = "3.24.0"
|
||||
openssl = "0.10.75"
|
||||
hmac = "0.12"
|
||||
sha2 = "0.10"
|
||||
hex = "0.4.3"
|
||||
flate2 = "1.1.5"
|
||||
tar = "0.4.44"
|
||||
|
||||
@@ -62,7 +62,7 @@ services:
|
||||
|
||||
db-mongodb-auth:
|
||||
container_name: db-mongodb-auth
|
||||
image: mongo:latest
|
||||
image: mongo:8.0.4
|
||||
ports:
|
||||
- "27082:27017"
|
||||
environment:
|
||||
@@ -75,14 +75,14 @@ services:
|
||||
volumes:
|
||||
- mongodb-data-auth:/data/db
|
||||
healthcheck:
|
||||
test: [ "CMD", "mongo", "--eval", "db.adminCommand('ping')" ]
|
||||
test: [ "CMD", "mongosh", "--eval", "db.adminCommand('ping')" ]
|
||||
interval: 5s
|
||||
timeout: 5s
|
||||
retries: 10
|
||||
|
||||
db-mongodb:
|
||||
container_name: db-mongodb
|
||||
image: mongo:latest
|
||||
image: mongo:8.0.4
|
||||
ports:
|
||||
- "27083:27017"
|
||||
volumes:
|
||||
|
||||
+3
-1
@@ -21,9 +21,11 @@ services:
|
||||
LOG: debug
|
||||
TZ: "Europe/Paris"
|
||||
# TMPDIR: /scratch
|
||||
EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiZDY4MzU2MTQtNzE2NC00OTQ4LWJlZjMtMTlkZDc5NGQzYmRhIiwibWFzdGVyS2V5QjY0IjoiV2NiM0pQQkVTaFBjRjg5UXZwRVJuamU4NGZmak1kNm4vS2dJOUpjMCtmVT0ifQ=="
|
||||
EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiYTQyNjQzNTQtZGE3Ni00OWFkLWJkYjctZDVjMjMwYzhjYmViIiwibWFzdGVyS2V5QjY0IjoiQlhWM1hvbEM2NTZTVjdkTmdjV1BHUWxrKytycExJNmxHRGk3Q1BCNWllbz0ifQ=="
|
||||
#CHUNK_SIZE_MB: "1"
|
||||
#POOLING: 1
|
||||
#RETRY_ATTEMPTS: 3
|
||||
#RETRY_BACKOFF_MS: 1000
|
||||
#DATABASES_CONFIG_FILE: "config.toml"
|
||||
extra_hosts:
|
||||
- "localhost:host-gateway"
|
||||
|
||||
+18
-1
@@ -9,7 +9,7 @@ RUN mkdir -p /mysql-exports/bin /mysql-exports/lib \
|
||||
# =========================
|
||||
# Base image (shared)
|
||||
# =========================
|
||||
FROM rust:1.94.0 AS base
|
||||
FROM rust:1.98 AS base
|
||||
|
||||
RUN apt-get update && DEBIAN_FRONTEND=noninteractive apt-get install -y \
|
||||
pkg-config \
|
||||
@@ -44,6 +44,15 @@ RUN ARCH=$(uname -m | sed 's/x86_64/amd64/;s/aarch64/arm64/') \
|
||||
| tar -xjf - -C /usr/local/bin sqlcmd \
|
||||
&& chmod +x /usr/local/bin/sqlcmd
|
||||
|
||||
# =========================
|
||||
# rclone (all storage backends)
|
||||
# =========================
|
||||
ARG RCLONE_VERSION=1.75.1
|
||||
RUN ARCH=$(dpkg --print-architecture) \
|
||||
&& curl -fsSL -o /tmp/rclone.deb "https://downloads.rclone.org/v${RCLONE_VERSION}/rclone-v${RCLONE_VERSION}-linux-${ARCH}.deb" \
|
||||
&& dpkg -i /tmp/rclone.deb \
|
||||
&& rm /tmp/rclone.deb
|
||||
|
||||
ARG TARGETARCH
|
||||
|
||||
# =========================
|
||||
@@ -134,6 +143,12 @@ RUN apt-get update && apt-get install -y \
|
||||
firebird3.0-utils \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
|
||||
ARG RCLONE_VERSION=1.75.1
|
||||
RUN ARCH=$(dpkg --print-architecture) \
|
||||
&& curl -fsSL -o /tmp/rclone.deb "https://downloads.rclone.org/v${RCLONE_VERSION}/rclone-v${RCLONE_VERSION}-linux-${ARCH}.deb" \
|
||||
&& dpkg -i /tmp/rclone.deb \
|
||||
&& rm /tmp/rclone.deb
|
||||
|
||||
ENV DOTNET_ROOT=/usr/local/dotnet
|
||||
RUN curl -sSL https://dot.net/v1/dotnet-install.sh -o /tmp/dotnet-install.sh \
|
||||
&& chmod +x /tmp/dotnet-install.sh \
|
||||
@@ -142,6 +157,8 @@ RUN curl -sSL https://dot.net/v1/dotnet-install.sh -o /tmp/dotnet-install.sh \
|
||||
|
||||
WORKDIR /app
|
||||
|
||||
RUN mkdir -p /config
|
||||
|
||||
COPY --from=builder /app/target/release/app /usr/local/bin/app
|
||||
COPY --from=builder /app/version.env /app/version.env
|
||||
COPY entrypoint.sh /entrypoint.sh
|
||||
|
||||
@@ -7,4 +7,6 @@ data:
|
||||
TZ: {{ .Values.env.TZ | quote }}
|
||||
POLLING: {{ .Values.env.POLLING | quote }}
|
||||
APP_ENV: {{ .Values.env.APP_ENV | quote }}
|
||||
LOG: {{ .Values.env.LOG | quote }}
|
||||
LOG: {{ .Values.env.LOG | quote }}
|
||||
RETRY_ATTEMPTS: {{ .Values.env.RETRY_ATTEMPTS | quote }}
|
||||
RETRY_BACKOFF_MS: {{ .Values.env.RETRY_BACKOFF_MS | quote }}
|
||||
|
||||
@@ -11,6 +11,8 @@ env:
|
||||
POLLING: "5"
|
||||
APP_ENV: "production"
|
||||
LOG: "info"
|
||||
RETRY_ATTEMPTS: "3"
|
||||
RETRY_BACKOFF_MS: "1000"
|
||||
|
||||
resources:
|
||||
limits:
|
||||
|
||||
@@ -20,6 +20,7 @@ pub async fn run(cfg: DatabaseConfig) -> anyhow::Result<bool> {
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.kill_on_drop(true)
|
||||
.spawn()?;
|
||||
|
||||
let query = b"SELECT 1 FROM RDB$DATABASE;\nQUIT;\n";
|
||||
|
||||
@@ -1,6 +1,5 @@
|
||||
use crate::domain::mariadb::connection::{select_mariadb_path, server_version};
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::domain::mysql::connection::connection_args;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use std::collections::HashMap;
|
||||
@@ -43,7 +42,9 @@ pub async fn run(
|
||||
|
||||
let start = Instant::now();
|
||||
let output = Command::new("mariadb-dump")
|
||||
.args(connection_args(&cfg))
|
||||
.arg("--host").arg(&cfg.host)
|
||||
.arg("--port").arg(cfg.port.to_string())
|
||||
.arg("--user").arg(&cfg.username)
|
||||
.arg("--routines")
|
||||
.arg("--events")
|
||||
.arg("--triggers")
|
||||
|
||||
@@ -1,12 +1,16 @@
|
||||
use std::path::PathBuf;
|
||||
use crate::domain::mysql::connection::connection_args;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::Result;
|
||||
use std::process::Command;
|
||||
|
||||
pub async fn server_version(cfg: &DatabaseConfig) -> Result<String> {
|
||||
let output = Command::new("mariadb")
|
||||
.args(connection_args(cfg))
|
||||
.arg("--host")
|
||||
.arg(&cfg.host)
|
||||
.arg("--port")
|
||||
.arg(cfg.port.to_string())
|
||||
.arg("--user")
|
||||
.arg(&cfg.username)
|
||||
.arg("-e")
|
||||
.arg("SELECT VERSION();")
|
||||
.env("MYSQL_PWD", &cfg.password)
|
||||
|
||||
@@ -1,4 +1,3 @@
|
||||
use crate::domain::mysql::connection::connection_args;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use std::collections::HashMap;
|
||||
use tokio::process::Command;
|
||||
@@ -6,9 +5,15 @@ use tokio::time::{Duration, timeout};
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, env: HashMap<String, String>) -> anyhow::Result<bool> {
|
||||
let mut cmd = Command::new("mysqladmin");
|
||||
cmd.args(connection_args(&cfg))
|
||||
cmd.arg("--host")
|
||||
.arg(cfg.host)
|
||||
.arg("--port")
|
||||
.arg(cfg.port.to_string())
|
||||
.arg("--user")
|
||||
.arg(cfg.username)
|
||||
.arg("ping")
|
||||
.envs(env);
|
||||
.envs(env)
|
||||
.kill_on_drop(true);
|
||||
|
||||
let result = timeout(Duration::from_secs(10), cmd.output()).await;
|
||||
|
||||
|
||||
@@ -1,5 +1,4 @@
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::domain::mysql::connection::connection_args;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use std::fs::File;
|
||||
@@ -22,7 +21,12 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogg
|
||||
|
||||
let drop_start = Instant::now();
|
||||
let drop_output = Command::new("mariadb")
|
||||
.args(connection_args(&cfg))
|
||||
.arg("--host")
|
||||
.arg(&cfg.host)
|
||||
.arg("--port")
|
||||
.arg(cfg.port.to_string())
|
||||
.arg("--user")
|
||||
.arg(&cfg.username)
|
||||
.arg("-e")
|
||||
.arg(&drop_create_cmd)
|
||||
.env("MYSQL_PWD", &cfg.password)
|
||||
@@ -44,7 +48,12 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogg
|
||||
let start = Instant::now();
|
||||
|
||||
let mut child = Command::new("mariadb")
|
||||
.args(connection_args(&cfg))
|
||||
.arg("--host")
|
||||
.arg(&cfg.host)
|
||||
.arg("--port")
|
||||
.arg(cfg.port.to_string())
|
||||
.arg("--user")
|
||||
.arg(&cfg.username)
|
||||
.arg("--database")
|
||||
.arg(&cfg.database)
|
||||
.env("MYSQL_PWD", &cfg.password)
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
pub mod docker_volume;
|
||||
pub mod factory;
|
||||
mod mongodb;
|
||||
pub mod mongodb;
|
||||
pub mod mysql;
|
||||
pub mod postgres;
|
||||
mod redis;
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::Result;
|
||||
use mongodb::Client;
|
||||
use percent_encoding::{utf8_percent_encode, AsciiSet, NON_ALPHANUMERIC};
|
||||
use percent_encoding::{AsciiSet, NON_ALPHANUMERIC, utf8_percent_encode};
|
||||
|
||||
const USERINFO_ENCODE: &AsciiSet = &NON_ALPHANUMERIC
|
||||
.remove(b'-')
|
||||
@@ -27,7 +27,8 @@ pub fn get_mongo_uri(cfg: DatabaseConfig) -> Result<String> {
|
||||
}
|
||||
|
||||
pub fn build_mongo_uri(cfg: &DatabaseConfig, include_db: bool) -> String {
|
||||
let is_srv = cfg.port == 0;
|
||||
let is_multi_host = cfg.host.contains(',');
|
||||
let is_srv = cfg.port == 0 && !is_multi_host;
|
||||
let scheme = if is_srv { "mongodb+srv" } else { "mongodb" };
|
||||
let has_auth = !cfg.username.is_empty() && !cfg.password.is_empty();
|
||||
|
||||
@@ -41,7 +42,7 @@ pub fn build_mongo_uri(cfg: &DatabaseConfig, include_db: bool) -> String {
|
||||
String::new()
|
||||
};
|
||||
|
||||
let authority = if is_srv {
|
||||
let authority = if is_srv || is_multi_host {
|
||||
cfg.host.clone()
|
||||
} else {
|
||||
format!("{}:{}", cfg.host, cfg.port)
|
||||
@@ -53,7 +54,35 @@ pub fn build_mongo_uri(cfg: &DatabaseConfig, include_db: bool) -> String {
|
||||
"/".to_string()
|
||||
};
|
||||
|
||||
let query = if has_auth { "?authSource=admin" } else { "" };
|
||||
let mut params: Vec<String> = Vec::new();
|
||||
|
||||
match cfg.options.get("auth_source").and_then(|v| v.as_str()) {
|
||||
Some(s) if !s.is_empty() => params.push(format!(
|
||||
"authSource={}",
|
||||
utf8_percent_encode(s, USERINFO_ENCODE)
|
||||
)),
|
||||
_ if has_auth => params.push("authSource=admin".to_string()),
|
||||
_ => {}
|
||||
}
|
||||
|
||||
if let Some(rs) = cfg.options.get("replica_set").and_then(|v| v.as_str()) {
|
||||
if !rs.is_empty() {
|
||||
params.push(format!(
|
||||
"replicaSet={}",
|
||||
utf8_percent_encode(rs, USERINFO_ENCODE)
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
if cfg.options.get("tls").and_then(|v| v.as_bool()) == Some(true) {
|
||||
params.push("tls=true".to_string());
|
||||
}
|
||||
|
||||
let query = if params.is_empty() {
|
||||
String::new()
|
||||
} else {
|
||||
format!("?{}", params.join("&"))
|
||||
};
|
||||
|
||||
format!("{}://{}{}{}{}", scheme, credentials, authority, path, query)
|
||||
}
|
||||
@@ -71,70 +100,3 @@ pub fn extract_db_name(dry_output: &str) -> Option<String> {
|
||||
}
|
||||
dbs.into_iter().next()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::services::config::{DatabaseConfig, DbType};
|
||||
use std::collections::HashMap;
|
||||
|
||||
fn cfg(host: &str, port: u16, user: &str, pass: &str) -> DatabaseConfig {
|
||||
DatabaseConfig {
|
||||
name: "t".into(),
|
||||
database: "mydb".into(),
|
||||
db_type: DbType::MongoDB,
|
||||
username: user.into(),
|
||||
password: pass.into(),
|
||||
port,
|
||||
host: host.into(),
|
||||
generated_id: "id".into(),
|
||||
path: String::new(),
|
||||
max_packet_size: String::new(),
|
||||
volume_name: String::new(),
|
||||
container_name: None,
|
||||
options: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_with_auth() {
|
||||
let c = cfg("localhost", 27017, "user", "pass");
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, true),
|
||||
"mongodb://user:pass@localhost:27017/mydb?authSource=admin"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_no_auth() {
|
||||
let c = cfg("localhost", 27017, "", "");
|
||||
assert_eq!(build_mongo_uri(&c, true), "mongodb://localhost:27017/mydb");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn srv_with_auth() {
|
||||
let c = cfg("cluster.example.mongodb.net", 0, "user", "pass");
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, true),
|
||||
"mongodb+srv://user:pass@cluster.example.mongodb.net/mydb?authSource=admin"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn srv_no_db_for_dryrun() {
|
||||
let c = cfg("cluster.example.mongodb.net", 0, "user", "pass");
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, false),
|
||||
"mongodb+srv://user:pass@cluster.example.mongodb.net/?authSource=admin"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn encodes_special_chars_in_credentials() {
|
||||
let c = cfg("cluster.example.mongodb.net", 0, "user", "p@ss:w/rd?");
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, true),
|
||||
"mongodb+srv://user:p%40ss%3Aw%2Frd%3F@cluster.example.mongodb.net/mydb?authSource=admin"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
mod backup;
|
||||
mod connection;
|
||||
pub mod connection;
|
||||
pub mod database;
|
||||
mod ping;
|
||||
mod restore;
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
use crate::domain::mysql::connection::{connection_args, server_version};
|
||||
use crate::domain::mysql::connection::server_version;
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
@@ -39,7 +39,9 @@ pub async fn run(
|
||||
|
||||
let start = Instant::now();
|
||||
let output = Command::new("mysqldump")
|
||||
.args(connection_args(&cfg))
|
||||
.arg("--host").arg(&cfg.host)
|
||||
.arg("--port").arg(cfg.port.to_string())
|
||||
.arg("--user").arg(&cfg.username)
|
||||
.arg("--routines")
|
||||
.arg("--events")
|
||||
.arg("--triggers")
|
||||
|
||||
@@ -2,33 +2,14 @@ use crate::services::config::DatabaseConfig;
|
||||
use anyhow::Result;
|
||||
use std::process::Command;
|
||||
|
||||
pub fn connection_args(cfg: &DatabaseConfig) -> Vec<String> {
|
||||
let protocol = cfg
|
||||
.options
|
||||
.get("protocol")
|
||||
.and_then(|v| v.as_str())
|
||||
.unwrap_or("tcp");
|
||||
|
||||
let mut args = vec![
|
||||
format!("--protocol={}", protocol),
|
||||
"--host".to_string(),
|
||||
cfg.host.clone(),
|
||||
"--port".to_string(),
|
||||
cfg.port.to_string(),
|
||||
"--user".to_string(),
|
||||
cfg.username.clone(),
|
||||
];
|
||||
|
||||
if let Some(socket) = cfg.options.get("socket").and_then(|v| v.as_str()) {
|
||||
args.push(format!("--socket={}", socket));
|
||||
}
|
||||
|
||||
args
|
||||
}
|
||||
|
||||
pub async fn server_version(cfg: &DatabaseConfig) -> Result<String> {
|
||||
let output = Command::new("mysql")
|
||||
.args(connection_args(cfg))
|
||||
.arg("--host")
|
||||
.arg(&cfg.host)
|
||||
.arg("--port")
|
||||
.arg(cfg.port.to_string())
|
||||
.arg("--user")
|
||||
.arg(&cfg.username)
|
||||
.arg("-e")
|
||||
.arg("SELECT VERSION();")
|
||||
.env("MYSQL_PWD", &cfg.password)
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
pub mod backup;
|
||||
pub mod connection;
|
||||
mod connection;
|
||||
pub mod database;
|
||||
mod ping;
|
||||
mod restore;
|
||||
|
||||
@@ -1,4 +1,3 @@
|
||||
use crate::domain::mysql::connection::connection_args;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use std::collections::HashMap;
|
||||
use tokio::process::Command;
|
||||
@@ -6,9 +5,15 @@ use tokio::time::{Duration, timeout};
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, env: HashMap<String, String>) -> anyhow::Result<bool> {
|
||||
let mut cmd = Command::new("mariadb-admin");
|
||||
cmd.args(connection_args(&cfg))
|
||||
cmd.arg("--host")
|
||||
.arg(cfg.host)
|
||||
.arg("--port")
|
||||
.arg(cfg.port.to_string())
|
||||
.arg("--user")
|
||||
.arg(cfg.username)
|
||||
.arg("ping")
|
||||
.envs(env);
|
||||
.envs(env)
|
||||
.kill_on_drop(true);
|
||||
|
||||
let result = timeout(Duration::from_secs(10), cmd.output()).await;
|
||||
|
||||
|
||||
@@ -1,5 +1,4 @@
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::domain::mysql::connection::connection_args;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use std::fs::File;
|
||||
@@ -22,7 +21,12 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogg
|
||||
|
||||
let drop_start = Instant::now();
|
||||
let drop_output = Command::new("mysql")
|
||||
.args(connection_args(&cfg))
|
||||
.arg("--host")
|
||||
.arg(&cfg.host)
|
||||
.arg("--port")
|
||||
.arg(cfg.port.to_string())
|
||||
.arg("--user")
|
||||
.arg(&cfg.username)
|
||||
.arg("-e")
|
||||
.arg(&drop_create_cmd)
|
||||
.env("MYSQL_PWD", &cfg.password)
|
||||
@@ -44,7 +48,12 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogg
|
||||
let start = Instant::now();
|
||||
|
||||
let mut child = Command::new("mysql")
|
||||
.args(connection_args(&cfg))
|
||||
.arg("--host")
|
||||
.arg(&cfg.host)
|
||||
.arg("--port")
|
||||
.arg(cfg.port.to_string())
|
||||
.arg("--user")
|
||||
.arg(&cfg.username)
|
||||
.arg("--database")
|
||||
.arg(&cfg.database)
|
||||
.env("MYSQL_PWD", &cfg.password)
|
||||
|
||||
@@ -21,6 +21,8 @@ pub async fn run(cfg: DatabaseConfig) -> Result<bool> {
|
||||
|
||||
cmd.arg("PING");
|
||||
|
||||
cmd.kill_on_drop(true);
|
||||
|
||||
debug!("Command Ping Redis: {:?}", cmd);
|
||||
|
||||
let result = timeout(Duration::from_secs(10), cmd.output()).await;
|
||||
|
||||
@@ -20,6 +20,7 @@ pub async fn run(cfg: DatabaseConfig) -> Result<bool> {
|
||||
}
|
||||
|
||||
cmd.arg("PING");
|
||||
cmd.kill_on_drop(true);
|
||||
|
||||
debug!("Command Ping Valkey: {:?}", cmd);
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
#![allow(dead_code)]
|
||||
|
||||
use crate::services::config::DbType;
|
||||
use std::fmt::{self, Display, Formatter};
|
||||
use std::path::PathBuf;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
@@ -20,3 +21,9 @@ pub struct UploadResult {
|
||||
pub remote_file_path: Option<String>,
|
||||
pub total_size: Option<u64>,
|
||||
}
|
||||
|
||||
impl Display for UploadResult {
|
||||
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
|
||||
write!(f, "{}", self.error.as_deref().unwrap_or("unknown error"))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ use super::service::BackupService;
|
||||
|
||||
use crate::domain::factory::DatabaseFactory;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use crate::utils::retry::{RetryPolicy, retry};
|
||||
|
||||
use anyhow::Result;
|
||||
use std::path::Path;
|
||||
@@ -39,7 +40,31 @@ impl BackupService {
|
||||
});
|
||||
}
|
||||
|
||||
match db.backup(tmp_path, Arc::clone(&logger)).await {
|
||||
let policy = RetryPolicy::default();
|
||||
|
||||
let db_ref = &db;
|
||||
let logger_ref = &logger;
|
||||
|
||||
let outcome = retry("Database backup", &logger, &policy, move |attempt| {
|
||||
let dir = tmp_path.join(format!("attempt-{attempt}"));
|
||||
|
||||
async move {
|
||||
if let Err(e) = tokio::fs::create_dir_all(&dir).await {
|
||||
return Err(anyhow::Error::from(e));
|
||||
}
|
||||
|
||||
match db_ref.backup(&dir, Arc::clone(logger_ref)).await {
|
||||
Ok(f) => Ok(f),
|
||||
Err(e) => {
|
||||
let _ = tokio::fs::remove_dir_all(&dir).await;
|
||||
Err(e)
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
.await;
|
||||
|
||||
match outcome {
|
||||
Ok(file) => Ok(BackupResult {
|
||||
generated_id,
|
||||
db_type,
|
||||
|
||||
@@ -4,6 +4,7 @@ use super::service::BackupService;
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::services::storage;
|
||||
use crate::utils::common::BackupMethod;
|
||||
use crate::utils::retry::{RetryPolicy, retry};
|
||||
use anyhow::{Result, bail};
|
||||
use futures::future::join_all;
|
||||
use std::sync::Arc;
|
||||
@@ -97,16 +98,44 @@ impl BackupService {
|
||||
/*
|
||||
STORAGE UPLOAD
|
||||
*/
|
||||
let upload_result = provider
|
||||
.upload(
|
||||
ctx_clone.clone(),
|
||||
result_clone,
|
||||
method,
|
||||
&storage,
|
||||
Some(encrypt),
|
||||
&backup_storage_id,
|
||||
let policy = RetryPolicy::default();
|
||||
|
||||
let attempt_result = if result_clone.backup_file.is_none() {
|
||||
logger_clone.log("error", format!("Missing backup file for storage {}", storage_id));
|
||||
|
||||
Err(UploadResult {
|
||||
storage_id: storage_id.clone(),
|
||||
success: false,
|
||||
error: Some("Missing backup file path".into()),
|
||||
remote_file_path: None,
|
||||
total_size: None,
|
||||
})
|
||||
} else {
|
||||
retry(
|
||||
&format!("Upload to storage {storage_id}"),
|
||||
&logger_clone,
|
||||
&policy,
|
||||
|_| async {
|
||||
let r = provider
|
||||
.upload(
|
||||
ctx_clone.clone(),
|
||||
result_clone.clone(),
|
||||
method,
|
||||
&storage,
|
||||
Some(encrypt),
|
||||
&backup_storage_id,
|
||||
)
|
||||
.await;
|
||||
|
||||
if r.success { Ok(r) } else { Err(r) }
|
||||
},
|
||||
)
|
||||
.await;
|
||||
.await
|
||||
};
|
||||
|
||||
let upload_result = match attempt_result {
|
||||
Ok(r) | Err(r) => r,
|
||||
};
|
||||
|
||||
let status = if upload_result.success { "success" } else { "failed" };
|
||||
|
||||
@@ -117,8 +146,6 @@ impl BackupService {
|
||||
upload_result.error.as_deref().unwrap_or("unknown error")
|
||||
));
|
||||
|
||||
// `backup_upload_init` opened a per-storage record; close it as "failed"
|
||||
// so the server is notified of the failure (no path/size on this path).
|
||||
if let Err(err) = ctx_clone
|
||||
.api
|
||||
.backup_upload_status(
|
||||
|
||||
+22
-7
@@ -207,16 +207,19 @@ impl ConfigService {
|
||||
ConfigService { ctx }
|
||||
}
|
||||
|
||||
pub fn load(&self, file_path: Option<&str>) -> Result<DatabasesConfig, String> {
|
||||
let path: String = if let Some(fp) = file_path {
|
||||
fp.to_string()
|
||||
} else {
|
||||
format!(
|
||||
fn resolve_path(file_path: Option<&str>) -> String {
|
||||
match file_path {
|
||||
Some(fp) => fp.to_string(),
|
||||
None => format!(
|
||||
"{}/{}",
|
||||
crate::settings::CONFIG.data_path,
|
||||
crate::settings::CONFIG.databases_config_file
|
||||
)
|
||||
};
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn load(&self, file_path: Option<&str>) -> Result<DatabasesConfig, String> {
|
||||
let path = Self::resolve_path(file_path);
|
||||
|
||||
info!("Loading databases config from: {}", path);
|
||||
|
||||
@@ -260,6 +263,18 @@ impl ConfigService {
|
||||
}
|
||||
|
||||
pub fn load_optional(&self, file_path: Option<&str>) -> DatabasesConfig {
|
||||
let path = Self::resolve_path(file_path);
|
||||
|
||||
if !Path::new(&path).exists() {
|
||||
info!(
|
||||
"No local databases config at {}; using dashboard-defined databases only",
|
||||
path
|
||||
);
|
||||
return DatabasesConfig {
|
||||
databases: Vec::new(),
|
||||
};
|
||||
}
|
||||
|
||||
self.load(file_path).unwrap_or_else(|e| {
|
||||
tracing::warn!(
|
||||
"Local databases config unavailable ({e}); continuing with dashboard-defined databases only"
|
||||
|
||||
@@ -46,6 +46,9 @@ pub fn persist_cache(path: &Path, databases: &[DatabaseConfig]) -> std::io::Resu
|
||||
};
|
||||
let json = serde_json::to_string_pretty(&wrapper)
|
||||
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
|
||||
if let Some(parent) = path.parent() {
|
||||
std::fs::create_dir_all(parent)?;
|
||||
}
|
||||
let tmp = path.with_extension("json.tmp");
|
||||
std::fs::write(&tmp, json)?;
|
||||
std::fs::rename(&tmp, path)?;
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
use super::service::RestoreService;
|
||||
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::utils::retry::{RetryPolicy, retry};
|
||||
use anyhow::Result;
|
||||
use futures::StreamExt;
|
||||
use reqwest::{Client, Url};
|
||||
@@ -7,7 +9,6 @@ use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
|
||||
fn human_size(bytes: u64) -> String {
|
||||
if bytes >= 1024 * 1024 {
|
||||
@@ -26,6 +27,34 @@ impl RestoreService {
|
||||
tmp_path: &Path,
|
||||
logger: Arc<JobLogger>,
|
||||
expected_size: Option<String>,
|
||||
) -> Result<PathBuf> {
|
||||
let policy = RetryPolicy::default();
|
||||
|
||||
let logger_ref = &logger;
|
||||
|
||||
let outcome = retry("Backup download", &logger, &policy, move |_| {
|
||||
let expected = expected_size.clone();
|
||||
|
||||
async move {
|
||||
self.download_once(file_url, tmp_path, Arc::clone(logger_ref), expected)
|
||||
.await
|
||||
}
|
||||
})
|
||||
.await;
|
||||
|
||||
if let Err(e) = &outcome {
|
||||
logger.log("error", format!("Download failed: {e}"));
|
||||
}
|
||||
|
||||
outcome
|
||||
}
|
||||
|
||||
pub async fn download_once(
|
||||
&self,
|
||||
file_url: &str,
|
||||
tmp_path: &Path,
|
||||
logger: Arc<JobLogger>,
|
||||
expected_size: Option<String>,
|
||||
) -> Result<PathBuf> {
|
||||
logger.log("info", "Start downloading backup archive".to_string());
|
||||
|
||||
@@ -69,7 +98,9 @@ impl RestoreService {
|
||||
format!(
|
||||
"Downloading backup '{}' ({})",
|
||||
filename,
|
||||
total.map(human_size).unwrap_or_else(|| "unknown size".to_string())
|
||||
total
|
||||
.map(human_size)
|
||||
.unwrap_or_else(|| "unknown size".to_string())
|
||||
),
|
||||
);
|
||||
|
||||
@@ -109,6 +140,16 @@ impl RestoreService {
|
||||
);
|
||||
}
|
||||
|
||||
if let Some(total) = total
|
||||
&& downloaded < total
|
||||
{
|
||||
anyhow::bail!(
|
||||
"Downloaded {} bytes but expected at least {} - backup appears truncated",
|
||||
downloaded,
|
||||
total
|
||||
);
|
||||
}
|
||||
|
||||
logger.log(
|
||||
"info",
|
||||
format!(
|
||||
|
||||
@@ -9,7 +9,9 @@ use providers::azure_blob;
|
||||
use providers::google_cloud_storage;
|
||||
use providers::google_drive;
|
||||
use providers::local;
|
||||
use providers::rclone;
|
||||
use providers::s3;
|
||||
use providers::sftp;
|
||||
use std::sync::Arc;
|
||||
use tracing::{error, info};
|
||||
|
||||
@@ -26,7 +28,6 @@ pub trait StorageProvider: Send + Sync {
|
||||
) -> UploadResult;
|
||||
}
|
||||
|
||||
/// Factory to create provider instance from storage config
|
||||
pub fn get_provider(storage: &DatabaseStorage) -> Option<Box<dyn StorageProvider>> {
|
||||
info!("Getting provider");
|
||||
info!("{:#?}", storage.provider.as_str());
|
||||
@@ -39,6 +40,8 @@ pub fn get_provider(storage: &DatabaseStorage) -> Option<Box<dyn StorageProvider
|
||||
"google-cloud-storage" => Some(Box::new(
|
||||
google_cloud_storage::GoogleCloudStorageProvider {},
|
||||
)),
|
||||
"rclone" => Some(Box::new(rclone::RcloneProvider {})),
|
||||
"sftp" => Some(Box::new(sftp::SftpProvider {})),
|
||||
_ => {
|
||||
error!("Unknown storage provider: {}", storage.provider);
|
||||
None
|
||||
|
||||
@@ -2,9 +2,8 @@ use anyhow::{Context as _, Result, anyhow};
|
||||
use base64::Engine;
|
||||
use base64::engine::general_purpose::STANDARD;
|
||||
use chrono::{Duration, Utc};
|
||||
use openssl::hash::MessageDigest;
|
||||
use openssl::pkey::PKey;
|
||||
use openssl::sign::Signer;
|
||||
use hmac::{Hmac, Mac};
|
||||
use sha2::Sha256;
|
||||
use url::Url;
|
||||
use azure_core::http::RequestContent;
|
||||
use azure_storage_blob::clients::{BlobClient, BlockBlobClient};
|
||||
@@ -36,12 +35,12 @@ impl SasResource {
|
||||
|
||||
pub(crate) const SAS_VERSION: &str = "2022-11-02";
|
||||
|
||||
type HmacSha256 = Hmac<Sha256>;
|
||||
|
||||
pub(crate) fn hmac_sha256_b64(key: &[u8], data: &str) -> Result<String> {
|
||||
let pkey = PKey::hmac(key).context("hmac key")?;
|
||||
let mut signer = Signer::new(MessageDigest::sha256(), &pkey).context("signer")?;
|
||||
signer.update(data.as_bytes()).context("signer update")?;
|
||||
let sig = signer.sign_to_vec().context("sign")?;
|
||||
Ok(STANDARD.encode(sig))
|
||||
let mut mac = HmacSha256::new_from_slice(key).context("hmac key")?;
|
||||
mac.update(data.as_bytes());
|
||||
Ok(STANDARD.encode(mac.finalize().into_bytes()))
|
||||
}
|
||||
|
||||
pub fn build_service_sas(
|
||||
|
||||
@@ -2,4 +2,6 @@ pub mod azure_blob;
|
||||
pub mod google_cloud_storage;
|
||||
pub mod google_drive;
|
||||
pub mod local;
|
||||
pub mod rclone;
|
||||
pub mod s3;
|
||||
pub mod sftp;
|
||||
|
||||
@@ -0,0 +1,205 @@
|
||||
use anyhow::{Context, Result, bail};
|
||||
use bytes::Bytes;
|
||||
use futures::{Stream, StreamExt};
|
||||
use std::io::Write;
|
||||
use std::path::Path;
|
||||
use std::pin::Pin;
|
||||
use std::process::Stdio;
|
||||
use tempfile::NamedTempFile;
|
||||
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||
use tokio::process::Command;
|
||||
use tracing::info;
|
||||
|
||||
const BLOCKED_BACKEND_TYPES: [&str; 13] = [
|
||||
"local",
|
||||
"alias",
|
||||
"crypt",
|
||||
"chunker",
|
||||
"compress",
|
||||
"union",
|
||||
"combine",
|
||||
"hasher",
|
||||
"archive",
|
||||
"cache",
|
||||
"memory",
|
||||
"http",
|
||||
"googlephotos",
|
||||
];
|
||||
|
||||
fn sections(config_text: &str) -> Vec<(String, Option<String>)> {
|
||||
let mut out: Vec<(String, Option<String>)> = Vec::new();
|
||||
|
||||
for line in config_text.lines() {
|
||||
let line = line.trim();
|
||||
|
||||
if line.starts_with('[') && line.ends_with(']') && line.len() > 2 {
|
||||
out.push((line[1..line.len() - 1].trim().to_string(), None));
|
||||
continue;
|
||||
}
|
||||
|
||||
let Some((key, value)) = line.split_once('=') else {
|
||||
continue;
|
||||
};
|
||||
|
||||
if key.trim().eq_ignore_ascii_case("type")
|
||||
&& let Some(current) = out.last_mut()
|
||||
&& current.1.is_none()
|
||||
{
|
||||
current.1 = Some(value.trim().to_ascii_lowercase());
|
||||
}
|
||||
}
|
||||
|
||||
out
|
||||
}
|
||||
|
||||
|
||||
pub fn validate_config(config_text: &str, remote_name: &str) -> Result<()> {
|
||||
let sections = sections(config_text);
|
||||
|
||||
if sections.is_empty() {
|
||||
bail!("rclone config contains no remote sections");
|
||||
}
|
||||
|
||||
for (name, backend) in §ions {
|
||||
let Some(backend) = backend else { continue };
|
||||
if BLOCKED_BACKEND_TYPES.contains(&backend.as_str()) {
|
||||
bail!("rclone backend type '{backend}' is not allowed (remote '{name}')");
|
||||
}
|
||||
}
|
||||
|
||||
if !sections.iter().any(|(name, _)| name == remote_name) {
|
||||
let available: Vec<&str> = sections.iter().map(|(name, _)| name.as_str()).collect();
|
||||
bail!(
|
||||
"remote '{remote_name}' is not defined in the rclone config (available: {})",
|
||||
available.join(", ")
|
||||
);
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn obscure_password(password: &str) -> Result<String> {
|
||||
let out = std::process::Command::new("rclone")
|
||||
.arg("obscure")
|
||||
.arg(password)
|
||||
.output()
|
||||
.context("failed to spawn rclone (is the binary installed in this image?)")?;
|
||||
|
||||
if !out.status.success() {
|
||||
bail!(
|
||||
"rclone obscure failed: {}",
|
||||
String::from_utf8_lossy(&out.stderr).trim()
|
||||
);
|
||||
}
|
||||
|
||||
Ok(String::from_utf8_lossy(&out.stdout).trim().to_string())
|
||||
}
|
||||
|
||||
pub fn build_rclone_config(remote_name: &str, fields: &[(&str, String)]) -> Result<String> {
|
||||
if remote_name.contains(['\r', '\n']) {
|
||||
bail!("rclone remote name must not contain line breaks");
|
||||
}
|
||||
|
||||
let mut lines = vec![format!("[{remote_name}]")];
|
||||
for (key, value) in fields {
|
||||
let value = value.trim();
|
||||
if value.is_empty() {
|
||||
continue;
|
||||
}
|
||||
if value.contains(['\r', '\n']) {
|
||||
bail!("rclone config value for '{key}' must not contain line breaks");
|
||||
}
|
||||
lines.push(format!("{key} = {value}"));
|
||||
}
|
||||
|
||||
Ok(lines.join("\n") + "\n")
|
||||
}
|
||||
|
||||
/// `<remote>:<remote_path>/<remote_file_path>`
|
||||
pub fn remote_target(remote_name: &str, remote_path: &str, remote_file_path: &str) -> String {
|
||||
let base = remote_path.trim().trim_matches('/');
|
||||
|
||||
if base.is_empty() {
|
||||
format!("{remote_name}:{remote_file_path}")
|
||||
} else {
|
||||
format!("{remote_name}:{base}/{remote_file_path}")
|
||||
}
|
||||
}
|
||||
|
||||
pub type RcloneStream = Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>>;
|
||||
|
||||
pub fn write_config(config_text: &str) -> Result<NamedTempFile> {
|
||||
let mut file = NamedTempFile::new().context("failed to create rclone config temp file")?;
|
||||
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
std::fs::set_permissions(file.path(), std::fs::Permissions::from_mode(0o600))
|
||||
.context("failed to restrict rclone config permissions")?;
|
||||
}
|
||||
|
||||
file.write_all(config_text.as_bytes())
|
||||
.context("failed to write rclone config")?;
|
||||
file.flush().context("failed to flush rclone config")?;
|
||||
|
||||
Ok(file)
|
||||
}
|
||||
|
||||
pub async fn rcat(config_path: &Path, target: &str, mut stream: RcloneStream) -> Result<()> {
|
||||
info!("rclone rcat -> {}", target);
|
||||
|
||||
let mut child = Command::new("rclone")
|
||||
.arg("--config")
|
||||
.arg(config_path)
|
||||
.arg("--contimeout")
|
||||
.arg("30s")
|
||||
.arg("--timeout")
|
||||
.arg("5m")
|
||||
.arg("--retries")
|
||||
.arg("1")
|
||||
.arg("--low-level-retries")
|
||||
.arg("3")
|
||||
.arg("rcat")
|
||||
.arg(target)
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::null())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()
|
||||
.context("failed to spawn rclone (is the binary installed in this image?)")?;
|
||||
|
||||
let mut stderr_pipe = child.stderr.take().context("rclone stderr unavailable")?;
|
||||
let stderr_task = tokio::spawn(async move {
|
||||
let mut buf = String::new();
|
||||
let _ = stderr_pipe.read_to_string(&mut buf).await;
|
||||
buf
|
||||
});
|
||||
|
||||
let mut stdin = child.stdin.take().context("rclone stdin unavailable")?;
|
||||
|
||||
while let Some(chunk) = stream.next().await {
|
||||
let chunk = match chunk {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
let _ = child.start_kill();
|
||||
let _ = child.wait().await;
|
||||
return Err(e).context("backup stream failed");
|
||||
}
|
||||
};
|
||||
|
||||
if stdin.write_all(&chunk).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
let _ = stdin.flush().await;
|
||||
drop(stdin);
|
||||
|
||||
let status = child.wait().await.context("failed to wait for rclone")?;
|
||||
let stderr = stderr_task.await.unwrap_or_default();
|
||||
|
||||
if !status.success() {
|
||||
bail!("rclone rcat failed ({status}): {}", stderr.trim());
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,112 @@
|
||||
pub mod helpers;
|
||||
pub mod models;
|
||||
|
||||
use crate::core::context::Context;
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::services::backup::models::{BackupResult, UploadResult};
|
||||
use crate::services::storage::StorageProvider;
|
||||
use crate::services::storage::providers::rclone::helpers::{
|
||||
rcat, remote_target, validate_config, write_config,
|
||||
};
|
||||
use crate::services::storage::providers::rclone::models::RcloneProviderConfig;
|
||||
use crate::utils::common::BackupMethod;
|
||||
use crate::utils::file::{full_file_name, full_file_path};
|
||||
use crate::utils::stream::build_stream;
|
||||
use async_trait::async_trait;
|
||||
use std::sync::Arc;
|
||||
use tokio::fs;
|
||||
use tracing::{error, info};
|
||||
|
||||
pub struct RcloneProvider {}
|
||||
|
||||
fn failed(storage_id: &str, error: impl ToString, total_size: Option<u64>) -> UploadResult {
|
||||
UploadResult {
|
||||
storage_id: storage_id.to_string(),
|
||||
success: false,
|
||||
error: Some(error.to_string()),
|
||||
remote_file_path: None,
|
||||
total_size,
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl StorageProvider for RcloneProvider {
|
||||
async fn upload(
|
||||
&self,
|
||||
ctx: Arc<Context>,
|
||||
result: BackupResult,
|
||||
_method: BackupMethod,
|
||||
storage: &DatabaseStorage,
|
||||
encrypt: Option<bool>,
|
||||
_backup_storage_id: &str,
|
||||
) -> UploadResult {
|
||||
let storage_id = storage.id.clone();
|
||||
|
||||
let Some(file_path) = result.backup_file else {
|
||||
return failed(&storage_id, "Missing backup file path", None);
|
||||
};
|
||||
|
||||
let total_size = match fs::metadata(&file_path).await {
|
||||
Ok(meta) => meta.len(),
|
||||
Err(e) => {
|
||||
error!("Failed to get file size: {}", e);
|
||||
return failed(&storage_id, e, None);
|
||||
}
|
||||
};
|
||||
|
||||
let config: RcloneProviderConfig = match storage.clone().config.try_into() {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
error!("rclone config deserialization failed: {}", e);
|
||||
return failed(&storage_id, e, Some(total_size));
|
||||
}
|
||||
};
|
||||
|
||||
if let Err(e) = validate_config(&config.config_text, &config.remote_name) {
|
||||
error!("rclone config rejected: {}", e);
|
||||
return failed(&storage_id, e, Some(total_size));
|
||||
}
|
||||
|
||||
let encrypt = encrypt.unwrap_or(false);
|
||||
|
||||
let upload = match build_stream(&file_path, encrypt, &ctx.edge_key.master_key_b64).await {
|
||||
Ok(u) => u,
|
||||
Err(e) => {
|
||||
error!("Stream build failed: {}", e);
|
||||
return failed(&storage_id, e, Some(total_size));
|
||||
}
|
||||
};
|
||||
|
||||
let file_name = full_file_name(encrypt);
|
||||
let remote_file_path = full_file_path(&file_name, storage.folder_name.as_deref());
|
||||
|
||||
let config_file = match write_config(&config.config_text) {
|
||||
Ok(f) => f,
|
||||
Err(e) => {
|
||||
error!("rclone config write failed: {}", e);
|
||||
return failed(&storage_id, e, Some(total_size));
|
||||
}
|
||||
};
|
||||
|
||||
let target = remote_target(&config.remote_name, &config.remote_path, &remote_file_path);
|
||||
|
||||
info!("Starting rclone upload to {}", target);
|
||||
|
||||
match rcat(config_file.path(), &target, upload.stream).await {
|
||||
Ok(()) => {
|
||||
info!("rclone upload successful: {}", remote_file_path);
|
||||
UploadResult {
|
||||
storage_id,
|
||||
success: true,
|
||||
error: None,
|
||||
remote_file_path: Some(remote_file_path),
|
||||
total_size: Some(total_size),
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
error!("rclone upload failed: {:?}", e);
|
||||
failed(&storage_id, e, Some(total_size))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
#[derive(Debug, Deserialize, Serialize)]
|
||||
pub struct RcloneProviderConfig {
|
||||
pub config_text: String,
|
||||
pub remote_name: String,
|
||||
pub remote_path: String,
|
||||
}
|
||||
@@ -0,0 +1,71 @@
|
||||
use crate::services::storage::providers::rclone::helpers::{build_rclone_config, obscure_password};
|
||||
use crate::services::storage::providers::sftp::models::SftpProviderConfig;
|
||||
use anyhow::{Context, Result, bail};
|
||||
use std::io::Write;
|
||||
use tempfile::NamedTempFile;
|
||||
|
||||
fn write_key(private_key: &str) -> Result<NamedTempFile> {
|
||||
let mut file = NamedTempFile::new().context("failed to create sftp key temp file")?;
|
||||
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
std::fs::set_permissions(file.path(), std::fs::Permissions::from_mode(0o600))
|
||||
.context("failed to restrict sftp key permissions")?;
|
||||
}
|
||||
|
||||
file.write_all(private_key.as_bytes())
|
||||
.context("failed to write sftp key")?;
|
||||
file.flush().context("failed to flush sftp key")?;
|
||||
Ok(file)
|
||||
}
|
||||
|
||||
pub fn build_sftp_config(
|
||||
config: &SftpProviderConfig,
|
||||
) -> Result<(String, Option<NamedTempFile>)> {
|
||||
if config.host.trim().is_empty() {
|
||||
bail!("sftp host is required");
|
||||
}
|
||||
if config.username.trim().is_empty() {
|
||||
bail!("sftp username is required");
|
||||
}
|
||||
|
||||
let has_password = config.password.as_deref().is_some_and(|p| !p.trim().is_empty());
|
||||
let has_key = config.private_key.as_deref().is_some_and(|k| !k.trim().is_empty());
|
||||
if !has_password && !has_key {
|
||||
bail!("sftp requires a password or a private key");
|
||||
}
|
||||
|
||||
let mut key_file: Option<NamedTempFile> = None;
|
||||
let mut key_file_path = String::new();
|
||||
if has_key {
|
||||
let file = write_key(config.private_key.as_deref().unwrap())?;
|
||||
key_file_path = file.path().display().to_string();
|
||||
key_file = Some(file);
|
||||
}
|
||||
|
||||
let pass = if has_password {
|
||||
obscure_password(config.password.as_deref().unwrap())?
|
||||
} else {
|
||||
String::new()
|
||||
};
|
||||
|
||||
let port = config
|
||||
.port
|
||||
.as_deref()
|
||||
.unwrap_or("")
|
||||
.trim()
|
||||
.to_string();
|
||||
|
||||
let fields: &[(&str, String)] = &[
|
||||
("type", "sftp".to_string()),
|
||||
("host", config.host.trim().to_string()),
|
||||
("port", port),
|
||||
("user", config.username.trim().to_string()),
|
||||
("key_file", key_file_path),
|
||||
("pass", pass),
|
||||
];
|
||||
|
||||
let config_text = build_rclone_config("sftp", fields)?;
|
||||
Ok((config_text, key_file))
|
||||
}
|
||||
@@ -0,0 +1,114 @@
|
||||
pub mod helpers;
|
||||
pub mod models;
|
||||
|
||||
use crate::core::context::Context;
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::services::backup::models::{BackupResult, UploadResult};
|
||||
use crate::services::storage::StorageProvider;
|
||||
use crate::services::storage::providers::rclone::helpers::{rcat, remote_target, write_config};
|
||||
use crate::services::storage::providers::sftp::helpers::build_sftp_config;
|
||||
use crate::services::storage::providers::sftp::models::SftpProviderConfig;
|
||||
use crate::utils::common::BackupMethod;
|
||||
use crate::utils::file::{full_file_name, full_file_path};
|
||||
use crate::utils::stream::build_stream;
|
||||
use async_trait::async_trait;
|
||||
use std::sync::Arc;
|
||||
use tokio::fs;
|
||||
use tracing::{error, info};
|
||||
|
||||
pub struct SftpProvider {}
|
||||
|
||||
fn failed(storage_id: &str, error: impl ToString, total_size: Option<u64>) -> UploadResult {
|
||||
UploadResult {
|
||||
storage_id: storage_id.to_string(),
|
||||
success: false,
|
||||
error: Some(error.to_string()),
|
||||
remote_file_path: None,
|
||||
total_size,
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl StorageProvider for SftpProvider {
|
||||
async fn upload(
|
||||
&self,
|
||||
ctx: Arc<Context>,
|
||||
result: BackupResult,
|
||||
_method: BackupMethod,
|
||||
storage: &DatabaseStorage,
|
||||
encrypt: Option<bool>,
|
||||
_backup_storage_id: &str,
|
||||
) -> UploadResult {
|
||||
let storage_id = storage.id.clone();
|
||||
|
||||
let Some(file_path) = result.backup_file else {
|
||||
return failed(&storage_id, "Missing backup file path", None);
|
||||
};
|
||||
|
||||
let total_size = match fs::metadata(&file_path).await {
|
||||
Ok(meta) => meta.len(),
|
||||
Err(e) => {
|
||||
error!("Failed to get file size: {}", e);
|
||||
return failed(&storage_id, e, None);
|
||||
}
|
||||
};
|
||||
|
||||
let config: SftpProviderConfig = match storage.clone().config.try_into() {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
error!("sftp config deserialization failed: {}", e);
|
||||
return failed(&storage_id, e, Some(total_size));
|
||||
}
|
||||
};
|
||||
|
||||
let (config_text, _key_file) = match build_sftp_config(&config) {
|
||||
Ok(v) => v,
|
||||
Err(e) => {
|
||||
error!("sftp config build failed: {}", e);
|
||||
return failed(&storage_id, e, Some(total_size));
|
||||
}
|
||||
};
|
||||
|
||||
let encrypt = encrypt.unwrap_or(false);
|
||||
|
||||
let upload = match build_stream(&file_path, encrypt, &ctx.edge_key.master_key_b64).await {
|
||||
Ok(u) => u,
|
||||
Err(e) => {
|
||||
error!("Stream build failed: {}", e);
|
||||
return failed(&storage_id, e, Some(total_size));
|
||||
}
|
||||
};
|
||||
|
||||
let file_name = full_file_name(encrypt);
|
||||
let remote_file_path = full_file_path(&file_name, storage.folder_name.as_deref());
|
||||
|
||||
let config_file = match write_config(&config_text) {
|
||||
Ok(f) => f,
|
||||
Err(e) => {
|
||||
error!("sftp config write failed: {}", e);
|
||||
return failed(&storage_id, e, Some(total_size));
|
||||
}
|
||||
};
|
||||
|
||||
let target = remote_target("sftp", &config.remote_path, &remote_file_path);
|
||||
|
||||
info!("Starting sftp (rclone) upload to {}", target);
|
||||
|
||||
match rcat(config_file.path(), &target, upload.stream).await {
|
||||
Ok(()) => {
|
||||
info!("sftp upload successful: {}", remote_file_path);
|
||||
UploadResult {
|
||||
storage_id,
|
||||
success: true,
|
||||
error: None,
|
||||
remote_file_path: Some(remote_file_path),
|
||||
total_size: Some(total_size),
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
error!("sftp upload failed: {:?}", e);
|
||||
failed(&storage_id, e, Some(total_size))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
use crate::utils::deserializer::string_or_number_to_string;
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
#[derive(Debug, Deserialize, Serialize)]
|
||||
pub struct SftpProviderConfig {
|
||||
pub host: String,
|
||||
#[serde(default, deserialize_with = "string_or_number_to_string")]
|
||||
pub port: Option<String>,
|
||||
pub username: String,
|
||||
#[serde(default)]
|
||||
pub password: Option<String>,
|
||||
#[serde(default)]
|
||||
pub private_key: Option<String>,
|
||||
#[serde(default)]
|
||||
pub remote_path: String,
|
||||
}
|
||||
+23
-1
@@ -16,6 +16,8 @@ pub struct Settings {
|
||||
pub timezone: String,
|
||||
pub log: String,
|
||||
pub chunk_size: usize, // bytes
|
||||
pub retry_attempts: u32,
|
||||
pub retry_backoff_ms: u64,
|
||||
}
|
||||
|
||||
impl Settings {
|
||||
@@ -49,6 +51,24 @@ impl Settings {
|
||||
|
||||
let chunk_size = chunk_size_mb * 1024 * 1024;
|
||||
|
||||
let retry_attempts = env::var("RETRY_ATTEMPTS")
|
||||
.unwrap_or_else(|_| "3".to_string())
|
||||
.parse::<u32>()
|
||||
.expect("RETRY_ATTEMPTS must be a valid positive integer");
|
||||
|
||||
if retry_attempts < 3 || retry_attempts > 5 {
|
||||
panic!("RETRY_ATTEMPTS must be between 3 and 5");
|
||||
}
|
||||
|
||||
let retry_backoff_ms = env::var("RETRY_BACKOFF_MS")
|
||||
.unwrap_or_else(|_| "1000".to_string())
|
||||
.parse::<u64>()
|
||||
.expect("RETRY_BACKOFF_MS must be a valid positive integer");
|
||||
|
||||
if retry_backoff_ms < 100 || retry_backoff_ms > 30_000 {
|
||||
panic!("RETRY_BACKOFF_MS must be between 100 and 30000 milliseconds");
|
||||
}
|
||||
|
||||
let tz = env::var("TZ").unwrap_or_else(|_| "UTC".to_string());
|
||||
|
||||
Self {
|
||||
@@ -64,7 +84,9 @@ impl Settings {
|
||||
pooling: pooling_seconds,
|
||||
timezone: tz,
|
||||
log: env::var("LOG").unwrap_or_else(|_| "info".into()),
|
||||
chunk_size
|
||||
chunk_size,
|
||||
retry_attempts,
|
||||
retry_backoff_ms,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -100,3 +100,102 @@ async fn mongodb_backup_restore_test() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
use crate::domain::mongodb::connection::build_mongo_uri;
|
||||
|
||||
fn uri_cfg(host: &str, port: u16, user: &str, pass: &str) -> DatabaseConfig {
|
||||
DatabaseConfig {
|
||||
name: "t".into(),
|
||||
database: "mydb".into(),
|
||||
db_type: DbType::MongoDB,
|
||||
username: user.into(),
|
||||
password: pass.into(),
|
||||
port,
|
||||
host: host.into(),
|
||||
generated_id: "id".into(),
|
||||
path: String::new(),
|
||||
max_packet_size: String::new(),
|
||||
volume_name: String::new(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn uri_standard_with_auth() {
|
||||
let c = uri_cfg("localhost", 27017, "user", "pass");
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, true),
|
||||
"mongodb://user:pass@localhost:27017/mydb?authSource=admin"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn uri_standard_no_auth() {
|
||||
let c = uri_cfg("localhost", 27017, "", "");
|
||||
assert_eq!(build_mongo_uri(&c, true), "mongodb://localhost:27017/mydb");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn uri_srv_with_auth() {
|
||||
let c = uri_cfg("cluster.example.mongodb.net", 0, "user", "pass");
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, true),
|
||||
"mongodb+srv://user:pass@cluster.example.mongodb.net/mydb?authSource=admin"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn uri_srv_no_db_for_dryrun() {
|
||||
let c = uri_cfg("cluster.example.mongodb.net", 0, "user", "pass");
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, false),
|
||||
"mongodb+srv://user:pass@cluster.example.mongodb.net/?authSource=admin"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn uri_options_authsource_replicaset_tls() {
|
||||
let mut c = uri_cfg("localhost", 27017, "user", "pass");
|
||||
c.options.insert("auth_source".into(), "myauthdb".into());
|
||||
c.options.insert("replica_set".into(), "rs0".into());
|
||||
c.options.insert("tls".into(), serde_json::Value::Bool(true));
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, true),
|
||||
"mongodb://user:pass@localhost:27017/mydb?authSource=myauthdb&replicaSet=rs0&tls=true"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn uri_multi_host_replica_set() {
|
||||
let mut c = uri_cfg(
|
||||
"mongodb0.example.internal:27017,mongodb1.example.internal:27017,mongodb2.example.internal:27017",
|
||||
0,
|
||||
"myDatabaseUser",
|
||||
"D1fficultP@ssw0rd",
|
||||
);
|
||||
c.database = "myDB".into();
|
||||
c.options.insert("replica_set".into(), "myRepl".into());
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, true),
|
||||
"mongodb://myDatabaseUser:D1fficultP%40ssw0rd@mongodb0.example.internal:27017,mongodb1.example.internal:27017,mongodb2.example.internal:27017/myDB?authSource=admin&replicaSet=myRepl"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn uri_default_authsource_when_auth() {
|
||||
let c = uri_cfg("localhost", 27017, "user", "pass");
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, true),
|
||||
"mongodb://user:pass@localhost:27017/mydb?authSource=admin"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn uri_encodes_special_chars_in_credentials() {
|
||||
let c = uri_cfg("cluster.example.mongodb.net", 0, "user", "p@ss:w/rd?");
|
||||
assert_eq!(
|
||||
build_mongo_uri(&c, true),
|
||||
"mongodb+srv://user:p%40ss%3Aw%2Frd%3F@cluster.example.mongodb.net/mydb?authSource=admin"
|
||||
);
|
||||
}
|
||||
|
||||
@@ -1,5 +1,4 @@
|
||||
use crate::domain::factory::DatabaseFactory;
|
||||
use crate::domain::mysql::connection::connection_args;
|
||||
use crate::services::config::{DatabaseConfig, DbType};
|
||||
use crate::tests::init_tracing_for_test;
|
||||
use crate::utils::compress::{compress_to_tar_gz_large, decompress_large_tar_gz};
|
||||
@@ -97,54 +96,3 @@ async fn mysql_backup_restore_test() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn tunnelled_config(options: serde_json::Value) -> DatabaseConfig {
|
||||
DatabaseConfig {
|
||||
name: "my-db".to_string(),
|
||||
database: "my-db".to_string(),
|
||||
db_type: DbType::Mysql,
|
||||
username: "my-db-user".to_string(),
|
||||
password: "my-db-password".to_string(),
|
||||
port: 3306,
|
||||
host: "localhost".to_string(),
|
||||
generated_id: "16678159-ff7e-4c97-8c83-0adeff214681".to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "512M".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: serde_json::from_value(options).unwrap(),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn connection_args_force_tcp_for_localhost() {
|
||||
// A `localhost` host makes the clients pick a Unix socket and ignore `--port`,
|
||||
// which breaks databases reached through an SSH tunnel.
|
||||
let cfg = tunnelled_config(serde_json::json!({}));
|
||||
|
||||
assert_eq!(
|
||||
connection_args(&cfg),
|
||||
vec![
|
||||
"--protocol=tcp",
|
||||
"--host",
|
||||
"localhost",
|
||||
"--port",
|
||||
"3306",
|
||||
"--user",
|
||||
"my-db-user",
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn connection_args_allow_socket_opt_in() {
|
||||
let cfg = tunnelled_config(serde_json::json!({
|
||||
"protocol": "socket",
|
||||
"socket": "/var/run/mysqld/mysqld.sock",
|
||||
}));
|
||||
|
||||
let args = connection_args(&cfg);
|
||||
|
||||
assert_eq!(args[0], "--protocol=socket");
|
||||
assert_eq!(args.last().unwrap(), "--socket=/var/run/mysqld/mysqld.sock");
|
||||
}
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
use crate::services::backup::BackupService;
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::{DatabaseConfig, DbType};
|
||||
use crate::tests::init_tracing_for_test;
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
use tempfile::TempDir;
|
||||
|
||||
fn sqlite_config(path: &str) -> DatabaseConfig {
|
||||
DatabaseConfig {
|
||||
name: "retry-test".to_string(),
|
||||
database: String::new(),
|
||||
db_type: DbType::Sqlite,
|
||||
username: String::new(),
|
||||
password: String::new(),
|
||||
port: 0,
|
||||
host: String::new(),
|
||||
generated_id: "retry-test-gen".to_string(),
|
||||
path: path.to_string(),
|
||||
max_packet_size: String::new(),
|
||||
volume_name: String::new(),
|
||||
container_name: None,
|
||||
options: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_failing_backup_is_retried_and_leaves_no_attempt_directory() {
|
||||
init_tracing_for_test();
|
||||
|
||||
let temp_dir = TempDir::new().unwrap();
|
||||
let tmp_path = temp_dir.path();
|
||||
let logger = Arc::new(JobLogger::new());
|
||||
|
||||
let cfg = sqlite_config("/nonexistent/definitely-not-here.sqlite");
|
||||
|
||||
let result = BackupService::run(cfg, tmp_path, Arc::clone(&logger))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(result.status, "failed");
|
||||
assert!(result.backup_file.is_none());
|
||||
|
||||
let entries = Arc::try_unwrap(logger).unwrap().into_entries();
|
||||
assert_eq!(
|
||||
entries.iter().filter(|e| e.level == "warn").count(),
|
||||
2,
|
||||
"expected one warn per non-final failed attempt"
|
||||
);
|
||||
assert!(
|
||||
entries
|
||||
.iter()
|
||||
.any(|e| e.level == "error" && e.message.starts_with("Backup failed:")),
|
||||
"expected a single terminal error from the runner"
|
||||
);
|
||||
|
||||
let leftovers: Vec<_> = std::fs::read_dir(tmp_path)
|
||||
.unwrap()
|
||||
.filter_map(|e| e.ok())
|
||||
.filter(|e| e.file_name().to_string_lossy().starts_with("attempt-"))
|
||||
.collect();
|
||||
assert!(
|
||||
leftovers.is_empty(),
|
||||
"failed attempt directories must be cleaned up, found {:?}",
|
||||
leftovers.iter().map(|e| e.file_name()).collect::<Vec<_>>()
|
||||
);
|
||||
}
|
||||
@@ -14,7 +14,9 @@ use crate::utils::common::BackupMethod;
|
||||
use crate::utils::edge_key::EdgeKey;
|
||||
|
||||
use serde_json::json;
|
||||
use std::io::Write;
|
||||
use std::sync::Arc;
|
||||
use tempfile::NamedTempFile;
|
||||
use wiremock::matchers::{body_partial_json, method, path};
|
||||
use wiremock::{Mock, MockServer, ResponseTemplate};
|
||||
|
||||
@@ -93,3 +95,111 @@ async fn failed_upload_reports_failed_status_to_server() {
|
||||
|
||||
// MockServer drop verifies both `.expect(1)` mounts were hit — including the "failed" PATCH.
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_failing_upload_is_retried_until_it_succeeds() {
|
||||
init_tracing_for_test();
|
||||
let server = MockServer::start().await;
|
||||
|
||||
Mock::given(method("POST"))
|
||||
.and(path("/agent/agent-1/backup/upload/init"))
|
||||
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
|
||||
"message": "ok",
|
||||
"backupStorage": { "id": "bs-1" }
|
||||
})))
|
||||
.expect(1)
|
||||
.mount(&server)
|
||||
.await;
|
||||
|
||||
Mock::given(method("POST"))
|
||||
.and(path("/tus/files"))
|
||||
.respond_with(ResponseTemplate::new(500))
|
||||
.up_to_n_times(2)
|
||||
.with_priority(1)
|
||||
.expect(2)
|
||||
.mount(&server)
|
||||
.await;
|
||||
|
||||
Mock::given(method("POST"))
|
||||
.and(path("/tus/files"))
|
||||
.respond_with(
|
||||
ResponseTemplate::new(201)
|
||||
.insert_header("Location", format!("{}/tus/files/upload-1", server.uri()).as_str()),
|
||||
)
|
||||
.with_priority(2)
|
||||
.expect(1)
|
||||
.mount(&server)
|
||||
.await;
|
||||
|
||||
Mock::given(method("PATCH"))
|
||||
.and(path("/tus/files/upload-1"))
|
||||
.respond_with(ResponseTemplate::new(204))
|
||||
.mount(&server)
|
||||
.await;
|
||||
|
||||
Mock::given(method("PATCH"))
|
||||
.and(path("/agent/agent-1/backup/upload/status"))
|
||||
.and(body_partial_json(json!({ "status": "success" })))
|
||||
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
|
||||
"message": "ok",
|
||||
"backupStorage": { "id": "bs-1" }
|
||||
})))
|
||||
.expect(1)
|
||||
.mount(&server)
|
||||
.await;
|
||||
|
||||
let mut backup_file = NamedTempFile::new().unwrap();
|
||||
backup_file.write_all(b"portabase-retry-test-payload").unwrap();
|
||||
backup_file.flush().unwrap();
|
||||
|
||||
let ctx = Context {
|
||||
edge_key: EdgeKey {
|
||||
server_url: server.uri(),
|
||||
agent_id: "agent-1".to_string(),
|
||||
master_key_b64: String::new(),
|
||||
},
|
||||
api: ApiClient::new(server.uri()),
|
||||
};
|
||||
|
||||
let service = BackupService::new(Arc::new(ctx));
|
||||
|
||||
let result = BackupResult {
|
||||
generated_id: "gen-1".to_string(),
|
||||
db_type: DbType::Postgresql,
|
||||
status: "success".to_string(),
|
||||
backup_file: Some(backup_file.path().to_path_buf()),
|
||||
code: None,
|
||||
};
|
||||
|
||||
let storage: DatabaseStorage = serde_json::from_value(json!({
|
||||
"id": "storage-1",
|
||||
"provider": "local",
|
||||
"config": {}
|
||||
}))
|
||||
.unwrap();
|
||||
|
||||
let backup_id = "backup-1".to_string();
|
||||
let logger = Arc::new(JobLogger::new());
|
||||
|
||||
let results = service
|
||||
.upload(
|
||||
result,
|
||||
BackupMethod::Manual,
|
||||
vec![storage],
|
||||
false,
|
||||
&backup_id,
|
||||
Arc::clone(&logger),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(results.len(), 1);
|
||||
assert!(results[0].success);
|
||||
|
||||
let entries = Arc::try_unwrap(logger).unwrap().into_entries();
|
||||
assert_eq!(entries.iter().filter(|e| e.level == "warn").count(), 2);
|
||||
assert!(
|
||||
entries.iter().any(|e| e.message
|
||||
== "Upload to storage storage-1 succeeded on attempt 3/3")
|
||||
);
|
||||
}
|
||||
|
||||
@@ -1,4 +1,6 @@
|
||||
mod api_models_tests;
|
||||
mod backup_runner_tests;
|
||||
mod backup_uploader_tests;
|
||||
mod config_tests;
|
||||
mod dashboard_config_tests;
|
||||
mod restore_downloader_tests;
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
use crate::core::context::Context;
|
||||
use crate::services::api::ApiClient;
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::restore::RestoreService;
|
||||
use crate::tests::init_tracing_for_test;
|
||||
use crate::utils::edge_key::EdgeKey;
|
||||
|
||||
use std::sync::Arc;
|
||||
use tempfile::TempDir;
|
||||
use wiremock::matchers::{method, path};
|
||||
use wiremock::{Mock, MockServer, ResponseTemplate};
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_failing_download_is_retried_until_it_succeeds() {
|
||||
init_tracing_for_test();
|
||||
let server = MockServer::start().await;
|
||||
|
||||
Mock::given(method("GET"))
|
||||
.and(path("/backups/archive.tar.gz"))
|
||||
.respond_with(ResponseTemplate::new(503))
|
||||
.up_to_n_times(2)
|
||||
.with_priority(1)
|
||||
.expect(2)
|
||||
.mount(&server)
|
||||
.await;
|
||||
|
||||
Mock::given(method("GET"))
|
||||
.and(path("/backups/archive.tar.gz"))
|
||||
.respond_with(ResponseTemplate::new(200).set_body_bytes(b"portabase-archive".to_vec()))
|
||||
.with_priority(2)
|
||||
.expect(1)
|
||||
.mount(&server)
|
||||
.await;
|
||||
|
||||
let ctx = Context {
|
||||
edge_key: EdgeKey {
|
||||
server_url: server.uri(),
|
||||
agent_id: "agent-1".to_string(),
|
||||
master_key_b64: String::new(),
|
||||
},
|
||||
api: ApiClient::new(server.uri()),
|
||||
};
|
||||
|
||||
let service = RestoreService::new(Arc::new(ctx));
|
||||
|
||||
let temp_dir = TempDir::new().unwrap();
|
||||
let logger = Arc::new(JobLogger::new());
|
||||
let url = format!("{}/backups/archive.tar.gz", server.uri());
|
||||
|
||||
let downloaded = service
|
||||
.download_backup(&url, temp_dir.path(), Arc::clone(&logger), None)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(std::fs::read(&downloaded).unwrap(), b"portabase-archive");
|
||||
|
||||
let entries = Arc::try_unwrap(logger).unwrap().into_entries();
|
||||
assert_eq!(entries.iter().filter(|e| e.level == "warn").count(), 2);
|
||||
assert!(
|
||||
entries
|
||||
.iter()
|
||||
.any(|e| e.message == "Backup download succeeded on attempt 3/3")
|
||||
);
|
||||
}
|
||||
@@ -1,2 +1,4 @@
|
||||
mod azure_blob;
|
||||
mod google_cloud_storage;
|
||||
mod rclone;
|
||||
mod sftp;
|
||||
|
||||
@@ -0,0 +1,455 @@
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::services::storage::providers::rclone::helpers::{remote_target, validate_config};
|
||||
use crate::services::storage::providers::rclone::models::RcloneProviderConfig;
|
||||
use crate::tests::init_tracing_for_test;
|
||||
use crate::utils::file::full_file_path;
|
||||
|
||||
const OVH_CONFIG: &str = "[ovhcloud-rbx]\n\
|
||||
type = s3\n\
|
||||
provider = OVHcloud\n\
|
||||
access_key_id = my_access\n\
|
||||
secret_access_key = my_secret\n\
|
||||
region = rbx\n\
|
||||
endpoint = s3.rbx.io.cloud.ovh.net\n\
|
||||
acl = private\n";
|
||||
|
||||
#[test]
|
||||
fn config_deserializes_from_dashboard_camel_case() {
|
||||
init_tracing_for_test();
|
||||
|
||||
let storage: DatabaseStorage = serde_json::from_value(serde_json::json!({
|
||||
"id": "storage-1",
|
||||
"provider": "rclone",
|
||||
"folderName": "backups",
|
||||
"config": {
|
||||
"configText": OVH_CONFIG,
|
||||
"remoteName": "ovhcloud-rbx",
|
||||
"remotePath": "my-bucket",
|
||||
}
|
||||
}))
|
||||
.unwrap();
|
||||
|
||||
let config: RcloneProviderConfig = storage.config.try_into().unwrap();
|
||||
|
||||
assert_eq!(config.remote_name, "ovhcloud-rbx");
|
||||
assert_eq!(config.remote_path, "my-bucket");
|
||||
assert!(config.config_text.contains("type = s3"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn build_rclone_config_skips_empty_and_rejects_line_breaks() {
|
||||
use crate::services::storage::providers::rclone::helpers::build_rclone_config;
|
||||
|
||||
let text = build_rclone_config(
|
||||
"sftp",
|
||||
&[
|
||||
("type", "sftp".to_string()),
|
||||
("host", "h".to_string()),
|
||||
("port", "".to_string()), // empty -> skipped
|
||||
("user", " u ".to_string()), // trimmed
|
||||
],
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(text, "[sftp]\ntype = sftp\nhost = h\nuser = u\n");
|
||||
|
||||
let err = build_rclone_config("sftp", &[("host", "a\nkey = injected".to_string())])
|
||||
.unwrap_err()
|
||||
.to_string();
|
||||
assert!(err.contains("line breaks"), "unexpected error: {err}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validate_config_accepts_the_target_remote() {
|
||||
assert!(validate_config(OVH_CONFIG, "ovhcloud-rbx").is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validate_config_rejects_an_unknown_remote_name() {
|
||||
let err = validate_config(OVH_CONFIG, "typo").unwrap_err().to_string();
|
||||
assert!(err.contains("typo"), "unexpected error: {err}");
|
||||
assert!(err.contains("ovhcloud-rbx"), "error should list the available remotes: {err}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validate_config_rejects_local_backend() {
|
||||
let cfg = "[disk]\ntype = local\n";
|
||||
let err = validate_config(cfg, "disk").unwrap_err().to_string();
|
||||
assert!(err.contains("local"), "unexpected error: {err}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validate_config_rejects_alias_backend() {
|
||||
let cfg = "[shortcut]\ntype = alias\nremote = other:path\n";
|
||||
let err = validate_config(cfg, "shortcut").unwrap_err().to_string();
|
||||
assert!(err.contains("alias"), "unexpected error: {err}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validate_config_rejects_a_blocked_backend_in_a_chained_section() {
|
||||
let cfg = "[secret]\ntype = crypt\nremote = disk:vault\n\n[disk]\ntype = local\n";
|
||||
let err = validate_config(cfg, "secret").unwrap_err().to_string();
|
||||
assert!(err.contains("crypt"), "unexpected error: {err}");
|
||||
assert!(err.contains("secret"), "error should name the offending remote: {err}");
|
||||
|
||||
let cfg = "[outer]\ntype = s3\nprovider = Minio\n\n[disk]\ntype = local\n";
|
||||
let err = validate_config(cfg, "outer").unwrap_err().to_string();
|
||||
assert!(err.contains("local"), "unexpected error: {err}");
|
||||
assert!(err.contains("disk"), "error should name the offending remote: {err}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validate_config_rejects_crypt_even_over_an_allowed_remote() {
|
||||
let cfg = format!("[secret]\ntype = crypt\nremote = ovhcloud-rbx:bucket\n\n{OVH_CONFIG}");
|
||||
let err = validate_config(&cfg, "secret").unwrap_err().to_string();
|
||||
assert!(err.contains("crypt"), "unexpected error: {err}");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validate_config_rejects_a_wrapping_backend_pointing_at_a_bare_local_path() {
|
||||
for backend in ["crypt", "chunker", "compress", "union", "combine", "hasher"] {
|
||||
let cfg = format!("[sneaky]\ntype = {backend}\nremote = /etc\n");
|
||||
let err = validate_config(&cfg, "sneaky")
|
||||
.unwrap_err()
|
||||
.to_string();
|
||||
assert!(err.contains(backend), "{backend} must be rejected: {err}");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validate_config_rejects_backends_that_cannot_hold_a_backup() {
|
||||
for backend in ["memory", "http", "googlephotos"] {
|
||||
let cfg = format!("[nope]\ntype = {backend}\n");
|
||||
let err = validate_config(&cfg, "nope").unwrap_err().to_string();
|
||||
assert!(err.contains(backend), "{backend} must be rejected: {err}");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn remote_path_is_a_prefix_ahead_of_the_backup_folder() {
|
||||
assert_eq!(
|
||||
remote_target("ovhcloud-rbx", "my-bucket", "backups/2026-09-09/x.tar.gz"),
|
||||
"ovhcloud-rbx:my-bucket/backups/2026-09-09/x.tar.gz"
|
||||
);
|
||||
|
||||
// Deeper prefixes nest the same way.
|
||||
assert_eq!(
|
||||
remote_target("ovhcloud-rbx", "my-bucket/portabase", "backups/2026-09-09/x.tar.gz"),
|
||||
"ovhcloud-rbx:my-bucket/portabase/backups/2026-09-09/x.tar.gz"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn remote_target_trims_surrounding_slashes_and_whitespace() {
|
||||
assert_eq!(
|
||||
remote_target("r", " /my-bucket/ ", "a/b.bin"),
|
||||
"r:my-bucket/a/b.bin"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn remote_target_handles_an_empty_remote_path() {
|
||||
assert_eq!(remote_target("r", "", "a/b.bin"), "r:a/b.bin");
|
||||
assert_eq!(remote_target("r", " ", "a/b.bin"), "r:a/b.bin");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_empty_remote_path_falls_back_to_the_global_backup_folder() {
|
||||
let remote_file_path = full_file_path(&"x.tar.gz".to_string(), None);
|
||||
assert!(remote_file_path.starts_with("backups/"));
|
||||
|
||||
assert_eq!(
|
||||
remote_target("ovhcloud-rbx", "", &remote_file_path),
|
||||
format!("ovhcloud-rbx:{remote_file_path}")
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
remote_target("ovhcloud-rbx", "my-bucket", &remote_file_path),
|
||||
format!("ovhcloud-rbx:my-bucket/{remote_file_path}")
|
||||
);
|
||||
}
|
||||
|
||||
use crate::services::storage::providers::rclone::helpers::{rcat, write_config};
|
||||
|
||||
use bytes::Bytes;
|
||||
use futures::stream;
|
||||
use std::process::Command;
|
||||
use testcontainers::core::{IntoContainerPort, WaitFor};
|
||||
use testcontainers::runners::AsyncRunner;
|
||||
use testcontainers::{GenericImage, ImageExt};
|
||||
|
||||
const BUCKET: &str = "portabase";
|
||||
|
||||
async fn start_minio() -> (testcontainers::ContainerAsync<GenericImage>, String) {
|
||||
let container = GenericImage::new("coollabsio/minio", "latest")
|
||||
.with_exposed_port(9000.tcp())
|
||||
.with_wait_for(WaitFor::message_on_stderr("API:"))
|
||||
.with_env_var("MINIO_ROOT_USER", "minioadmin")
|
||||
.with_env_var("MINIO_ROOT_PASSWORD", "minioadmin")
|
||||
.with_cmd(["server", "/data"])
|
||||
.start()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let host = container.get_host().await.unwrap().to_string();
|
||||
let port = container.get_host_port_ipv4(9000).await.unwrap();
|
||||
(container, format!("http://{host}:{port}"))
|
||||
}
|
||||
|
||||
fn minio_config(endpoint: &str) -> String {
|
||||
format!(
|
||||
"[minio]\n\
|
||||
type = s3\n\
|
||||
provider = Minio\n\
|
||||
access_key_id = minioadmin\n\
|
||||
secret_access_key = minioadmin\n\
|
||||
endpoint = {endpoint}\n\
|
||||
region = us-east-1\n\
|
||||
force_path_style = true\n"
|
||||
)
|
||||
}
|
||||
|
||||
fn rclone_ok(config_path: &std::path::Path, args: &[&str]) -> Vec<u8> {
|
||||
let out = Command::new("rclone")
|
||||
.arg("--config")
|
||||
.arg(config_path)
|
||||
.args(args)
|
||||
.output()
|
||||
.expect("rclone binary not found — is it installed in this image?");
|
||||
|
||||
assert!(
|
||||
out.status.success(),
|
||||
"rclone {args:?} failed: {}",
|
||||
String::from_utf8_lossy(&out.stderr)
|
||||
);
|
||||
|
||||
out.stdout
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_config_creates_an_owner_only_file_with_the_exact_text() {
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
|
||||
let file = write_config(OVH_CONFIG).unwrap();
|
||||
|
||||
let mode = std::fs::metadata(file.path()).unwrap().permissions().mode();
|
||||
assert_eq!(mode & 0o777, 0o600, "config file must not be group/world readable");
|
||||
|
||||
assert_eq!(std::fs::read_to_string(file.path()).unwrap(), OVH_CONFIG);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rcat_streams_a_multi_chunk_body_to_minio() {
|
||||
init_tracing_for_test();
|
||||
|
||||
let (_container, endpoint) = start_minio().await;
|
||||
let config = write_config(&minio_config(&endpoint)).unwrap();
|
||||
|
||||
rclone_ok(config.path(), &["mkdir", &format!("minio:{BUCKET}")]);
|
||||
|
||||
let data = vec![7u8; 10 * 1024];
|
||||
let chunks: Vec<Result<Bytes, std::io::Error>> = data
|
||||
.chunks(1024)
|
||||
.map(|c| Ok(Bytes::copy_from_slice(c)))
|
||||
.collect();
|
||||
|
||||
let target = remote_target("minio", BUCKET, "backups/2026-09-09/test.bin");
|
||||
|
||||
rcat(config.path(), &target, Box::pin(stream::iter(chunks)))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let got = rclone_ok(config.path(), &["cat", &target]);
|
||||
assert_eq!(got, data);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rcat_reports_rclone_stderr_when_the_remote_is_unreachable() {
|
||||
init_tracing_for_test();
|
||||
|
||||
let config = write_config(&minio_config("http://127.0.0.1:1")).unwrap();
|
||||
|
||||
let chunks: Vec<Result<Bytes, std::io::Error>> =
|
||||
vec![Ok(Bytes::from_static(&[0u8; 4096]))];
|
||||
|
||||
let err = rcat(
|
||||
config.path(),
|
||||
"minio:portabase/x.bin",
|
||||
Box::pin(stream::iter(chunks)),
|
||||
)
|
||||
.await
|
||||
.expect_err("upload to an unreachable endpoint must fail");
|
||||
|
||||
let msg = err.to_string();
|
||||
assert!(
|
||||
msg.contains("rclone rcat failed"),
|
||||
"the broken stdin pipe must not mask rclone's own error: {msg}"
|
||||
);
|
||||
let (_, stderr_part) = msg
|
||||
.rsplit_once(": ")
|
||||
.expect("bail message must carry rclone stderr after the exit status");
|
||||
assert!(
|
||||
!stderr_part.trim().is_empty(),
|
||||
"rclone stderr must be included: {msg}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn rcat_aborts_the_upload_when_the_stream_fails() {
|
||||
init_tracing_for_test();
|
||||
|
||||
let (_container, endpoint) = start_minio().await;
|
||||
let config = write_config(&minio_config(&endpoint)).unwrap();
|
||||
|
||||
rclone_ok(config.path(), &["mkdir", &format!("minio:{BUCKET}")]);
|
||||
|
||||
let chunks: Vec<Result<Bytes, std::io::Error>> = vec![
|
||||
Ok(Bytes::from_static(&[1u8; 1024])),
|
||||
Err(std::io::Error::other("injected stream failure")),
|
||||
];
|
||||
|
||||
let target = remote_target("minio", BUCKET, "backups/2026-09-09/aborted.bin");
|
||||
|
||||
let err = rcat(config.path(), &target, Box::pin(stream::iter(chunks)))
|
||||
.await
|
||||
.expect_err("a stream error must fail the upload");
|
||||
assert!(
|
||||
err.to_string().contains("backup stream failed"),
|
||||
"unexpected error: {err}"
|
||||
);
|
||||
|
||||
let stat_out = rclone_ok(config.path(), &["lsjson", "--stat", &target]);
|
||||
let stat: serde_json::Value = serde_json::from_slice(&stat_out).unwrap();
|
||||
assert_eq!(
|
||||
stat["Name"], "",
|
||||
"rclone must not have finalized the truncated object: {stat}"
|
||||
);
|
||||
assert_eq!(stat["IsDir"], true, "a miss reports IsDir: true: {stat}");
|
||||
}
|
||||
|
||||
use crate::core::context::Context;
|
||||
use crate::services::api::ApiClient;
|
||||
use crate::services::backup::models::BackupResult;
|
||||
use crate::services::config::DbType;
|
||||
use crate::services::storage::providers::rclone::RcloneProvider;
|
||||
use crate::services::storage::{StorageProvider, get_provider};
|
||||
use crate::utils::common::BackupMethod;
|
||||
use crate::utils::edge_key::EdgeKey;
|
||||
|
||||
use std::io::Write as _;
|
||||
use std::sync::Arc;
|
||||
use tempfile::NamedTempFile;
|
||||
|
||||
fn test_context() -> Arc<Context> {
|
||||
Arc::new(Context {
|
||||
edge_key: EdgeKey {
|
||||
server_url: String::new(),
|
||||
agent_id: "agent-1".to_string(),
|
||||
master_key_b64: String::new(),
|
||||
},
|
||||
api: ApiClient::new(String::new()),
|
||||
})
|
||||
}
|
||||
|
||||
fn storage_for(config_text: &str, remote_path: &str) -> DatabaseStorage {
|
||||
serde_json::from_value(serde_json::json!({
|
||||
"id": "storage-1",
|
||||
"provider": "rclone",
|
||||
"folderName": "backups",
|
||||
"config": {
|
||||
"configText": config_text,
|
||||
"remoteName": "minio",
|
||||
"remotePath": remote_path,
|
||||
}
|
||||
}))
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn factory_resolves_the_rclone_provider_key() {
|
||||
let storage = storage_for(OVH_CONFIG, "bucket");
|
||||
assert!(
|
||||
get_provider(&storage).is_some(),
|
||||
"get_provider must recognise the \"rclone\" key"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn provider_uploads_an_unencrypted_backup_to_minio() {
|
||||
init_tracing_for_test();
|
||||
|
||||
let (_container, endpoint) = start_minio().await;
|
||||
let config_text = minio_config(&endpoint);
|
||||
|
||||
let bootstrap = write_config(&config_text).unwrap();
|
||||
rclone_ok(bootstrap.path(), &["mkdir", &format!("minio:{BUCKET}")]);
|
||||
|
||||
let payload = vec![42u8; 64 * 1024];
|
||||
let mut backup_file = NamedTempFile::new().unwrap();
|
||||
backup_file.write_all(&payload).unwrap();
|
||||
backup_file.flush().unwrap();
|
||||
|
||||
let storage = storage_for(&config_text, BUCKET);
|
||||
|
||||
let result = RcloneProvider {}
|
||||
.upload(
|
||||
test_context(),
|
||||
BackupResult {
|
||||
generated_id: "db-1".to_string(),
|
||||
db_type: DbType::Postgresql,
|
||||
status: "success".to_string(),
|
||||
backup_file: Some(backup_file.path().to_path_buf()),
|
||||
code: None,
|
||||
},
|
||||
BackupMethod::Automatic,
|
||||
&storage,
|
||||
Some(false),
|
||||
"backup-storage-1",
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(result.success, "upload failed: {:?}", result.error);
|
||||
assert_eq!(result.total_size, Some(payload.len() as u64));
|
||||
|
||||
let remote_file_path = result.remote_file_path.expect("remote path must be reported");
|
||||
assert!(
|
||||
remote_file_path.starts_with("backups/"),
|
||||
"folder_name must prefix the path: {remote_file_path}"
|
||||
);
|
||||
|
||||
let target = remote_target("minio", BUCKET, &remote_file_path);
|
||||
assert_eq!(rclone_ok(bootstrap.path(), &["cat", &target]), payload);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn provider_refuses_a_blocked_backend_without_spawning_rclone() {
|
||||
init_tracing_for_test();
|
||||
|
||||
let mut backup_file = NamedTempFile::new().unwrap();
|
||||
backup_file.write_all(b"payload").unwrap();
|
||||
backup_file.flush().unwrap();
|
||||
|
||||
let storage = storage_for("[minio]\ntype = local\n", "bucket");
|
||||
|
||||
let result = RcloneProvider {}
|
||||
.upload(
|
||||
test_context(),
|
||||
BackupResult {
|
||||
generated_id: "db-1".to_string(),
|
||||
db_type: DbType::Postgresql,
|
||||
status: "success".to_string(),
|
||||
backup_file: Some(backup_file.path().to_path_buf()),
|
||||
code: None,
|
||||
},
|
||||
BackupMethod::Automatic,
|
||||
&storage,
|
||||
Some(false),
|
||||
"backup-storage-1",
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(!result.success);
|
||||
assert!(
|
||||
result.error.unwrap_or_default().contains("local"),
|
||||
"the error must name the rejected backend type"
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::services::storage::providers::sftp::helpers::build_sftp_config;
|
||||
use crate::services::storage::providers::sftp::models::SftpProviderConfig;
|
||||
use crate::services::storage::get_provider;
|
||||
|
||||
fn storage(config: serde_json::Value) -> DatabaseStorage {
|
||||
serde_json::from_value(serde_json::json!({
|
||||
"id": "storage-1",
|
||||
"provider": "sftp",
|
||||
"folderName": "backups",
|
||||
"config": config,
|
||||
}))
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn config_deserializes_from_dashboard_camel_case() {
|
||||
let s = storage(serde_json::json!({
|
||||
"host": "backup.example.com",
|
||||
"port": 2222,
|
||||
"username": "deploy",
|
||||
"privateKey": "-----BEGIN KEY-----",
|
||||
"remotePath": "/srv/backups",
|
||||
}));
|
||||
let config: SftpProviderConfig = s.config.try_into().unwrap();
|
||||
assert_eq!(config.host, "backup.example.com");
|
||||
assert_eq!(config.port.as_deref(), Some("2222"));
|
||||
assert_eq!(config.username, "deploy");
|
||||
assert_eq!(config.remote_path, "/srv/backups");
|
||||
assert_eq!(config.private_key.as_deref(), Some("-----BEGIN KEY-----"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn build_config_emits_key_file_for_key_auth() {
|
||||
let config = SftpProviderConfig {
|
||||
host: "h".into(),
|
||||
port: Some("2222".into()),
|
||||
username: "u".into(),
|
||||
password: None,
|
||||
private_key: Some("PEMDATA".into()),
|
||||
remote_path: String::new(),
|
||||
};
|
||||
let (text, key) = build_sftp_config(&config).unwrap();
|
||||
let key = key.expect("key auth must produce a key file");
|
||||
assert!(text.contains("[sftp]"));
|
||||
assert!(text.contains("type = sftp"));
|
||||
assert!(text.contains("host = h"));
|
||||
assert!(text.contains("port = 2222"));
|
||||
assert!(text.contains("user = u"));
|
||||
assert!(text.contains(&format!("key_file = {}", key.path().display())));
|
||||
assert_eq!(std::fs::read_to_string(key.path()).unwrap(), "PEMDATA");
|
||||
assert!(!text.contains("pass ="));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn build_config_omits_port_when_absent() {
|
||||
let config = SftpProviderConfig {
|
||||
host: "h".into(),
|
||||
port: None,
|
||||
username: "u".into(),
|
||||
password: None,
|
||||
private_key: Some("K".into()),
|
||||
remote_path: String::new(),
|
||||
};
|
||||
let (text, _key) = build_sftp_config(&config).unwrap();
|
||||
assert!(!text.contains("port ="));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn factory_resolves_the_sftp_provider_key() {
|
||||
let s = storage(serde_json::json!({
|
||||
"host": "h", "username": "u", "password": "pw",
|
||||
}));
|
||||
assert!(get_provider(&s).is_some());
|
||||
}
|
||||
@@ -4,4 +4,5 @@ mod deserializer;
|
||||
mod edge_key_tests;
|
||||
mod file_tests;
|
||||
mod normalize_cron_tests;
|
||||
mod retry_tests;
|
||||
mod stream_tests;
|
||||
|
||||
@@ -0,0 +1,135 @@
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::tests::init_tracing_for_test;
|
||||
use crate::utils::retry::{RetryPolicy, retry};
|
||||
|
||||
use std::sync::Mutex;
|
||||
use std::sync::atomic::{AtomicU32, Ordering};
|
||||
use std::time::Duration;
|
||||
|
||||
fn fast_policy(attempts: u32) -> RetryPolicy {
|
||||
RetryPolicy {
|
||||
attempts,
|
||||
base_backoff: Duration::from_millis(1),
|
||||
max_backoff: Duration::from_millis(4),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn delay_grows_with_the_attempt_number() {
|
||||
let policy = RetryPolicy {
|
||||
attempts: 5,
|
||||
base_backoff: Duration::from_millis(100),
|
||||
max_backoff: Duration::from_secs(30),
|
||||
};
|
||||
|
||||
assert!(policy.delay(1) >= Duration::from_millis(50));
|
||||
assert!(policy.delay(1) <= Duration::from_millis(100));
|
||||
assert!(policy.delay(2) >= Duration::from_millis(100));
|
||||
assert!(policy.delay(2) <= Duration::from_millis(200));
|
||||
assert!(policy.delay(3) >= Duration::from_millis(200));
|
||||
assert!(policy.delay(3) <= Duration::from_millis(400));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn delay_never_exceeds_max_backoff() {
|
||||
let policy = RetryPolicy {
|
||||
attempts: 5,
|
||||
base_backoff: Duration::from_millis(1000),
|
||||
max_backoff: Duration::from_millis(2000),
|
||||
};
|
||||
|
||||
for attempt in 1..=5 {
|
||||
assert!(policy.delay(attempt) <= Duration::from_millis(2000));
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn first_attempt_success_logs_nothing() {
|
||||
init_tracing_for_test();
|
||||
let logger = JobLogger::new();
|
||||
|
||||
let result: Result<u32, anyhow::Error> =
|
||||
retry("Test op", &logger, &fast_policy(3), |_| async { Ok(7) }).await;
|
||||
|
||||
assert_eq!(result.unwrap(), 7);
|
||||
assert!(logger.into_entries().is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn retries_until_success_and_logs_each_attempt() {
|
||||
init_tracing_for_test();
|
||||
let logger = JobLogger::new();
|
||||
let calls = AtomicU32::new(0);
|
||||
|
||||
let result: Result<u32, anyhow::Error> =
|
||||
retry("Test op", &logger, &fast_policy(3), |_| async {
|
||||
let n = calls.fetch_add(1, Ordering::SeqCst) + 1;
|
||||
if n < 3 {
|
||||
Err(anyhow::anyhow!("boom {n}"))
|
||||
} else {
|
||||
Ok(n)
|
||||
}
|
||||
})
|
||||
.await;
|
||||
|
||||
assert_eq!(result.unwrap(), 3);
|
||||
assert_eq!(calls.load(Ordering::SeqCst), 3);
|
||||
|
||||
let entries = logger.into_entries();
|
||||
|
||||
let warns: Vec<_> = entries.iter().filter(|e| e.level == "warn").collect();
|
||||
assert_eq!(warns.len(), 2);
|
||||
assert!(warns[0].message.starts_with("Test op attempt 1/3 failed: boom 1"));
|
||||
assert!(warns[1].message.starts_with("Test op attempt 2/3 failed: boom 2"));
|
||||
|
||||
let infos: Vec<_> = entries.iter().filter(|e| e.level == "info").collect();
|
||||
assert_eq!(infos.len(), 1);
|
||||
assert_eq!(infos[0].message, "Test op succeeded on attempt 3/3");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn exhausts_attempts_and_logs_no_terminal_error() {
|
||||
init_tracing_for_test();
|
||||
let logger = JobLogger::new();
|
||||
let calls = AtomicU32::new(0);
|
||||
|
||||
let result: Result<(), anyhow::Error> =
|
||||
retry("Test op", &logger, &fast_policy(3), |_| async {
|
||||
calls.fetch_add(1, Ordering::SeqCst);
|
||||
Err(anyhow::anyhow!("always"))
|
||||
})
|
||||
.await;
|
||||
|
||||
assert!(result.is_err());
|
||||
assert_eq!(calls.load(Ordering::SeqCst), 3);
|
||||
|
||||
let entries = logger.into_entries();
|
||||
assert_eq!(entries.iter().filter(|e| e.level == "warn").count(), 2);
|
||||
assert_eq!(entries.iter().filter(|e| e.level == "error").count(), 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn closure_receives_the_attempt_number() {
|
||||
init_tracing_for_test();
|
||||
let logger = JobLogger::new();
|
||||
let seen = Mutex::new(Vec::new());
|
||||
let seen_ref = &seen;
|
||||
|
||||
let result: Result<(), anyhow::Error> =
|
||||
retry("Test op", &logger, &fast_policy(3), move |attempt| async move {
|
||||
seen_ref.lock().unwrap().push(attempt);
|
||||
Err(anyhow::anyhow!("nope"))
|
||||
})
|
||||
.await;
|
||||
|
||||
assert!(result.is_err());
|
||||
assert_eq!(*seen.lock().unwrap(), vec![1, 2, 3]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn config_defaults_are_within_the_documented_range() {
|
||||
let policy = RetryPolicy::default();
|
||||
assert!(policy.attempts >= 3 && policy.attempts <= 5);
|
||||
assert!(policy.base_backoff >= Duration::from_millis(100));
|
||||
assert!(policy.base_backoff <= Duration::from_millis(30_000));
|
||||
}
|
||||
@@ -6,6 +6,7 @@ pub mod file;
|
||||
pub mod locks;
|
||||
pub mod logging;
|
||||
pub mod redis_client;
|
||||
pub mod retry;
|
||||
pub mod stream;
|
||||
pub mod task_manager;
|
||||
pub mod text;
|
||||
|
||||
@@ -0,0 +1,75 @@
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::settings::CONFIG;
|
||||
use rand::Rng;
|
||||
use std::fmt::Display;
|
||||
use std::future::Future;
|
||||
use std::time::Duration;
|
||||
|
||||
pub struct RetryPolicy {
|
||||
pub attempts: u32,
|
||||
pub base_backoff: Duration,
|
||||
pub max_backoff: Duration,
|
||||
}
|
||||
|
||||
impl Default for RetryPolicy {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
attempts: CONFIG.retry_attempts,
|
||||
base_backoff: Duration::from_millis(CONFIG.retry_backoff_ms),
|
||||
max_backoff: Duration::from_secs(30),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl RetryPolicy {
|
||||
pub(crate) fn delay(&self, attempt: u32) -> Duration {
|
||||
let exp = self.base_backoff.saturating_mul(1u32 << (attempt - 1).min(16));
|
||||
let capped = exp.min(self.max_backoff);
|
||||
let half = capped / 2;
|
||||
let jitter = rand::rng().random_range(0..=half.as_millis() as u64);
|
||||
|
||||
half + Duration::from_millis(jitter)
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn retry<T, E, F, Fut>(
|
||||
op: &str,
|
||||
logger: &JobLogger,
|
||||
policy: &RetryPolicy,
|
||||
mut f: F,
|
||||
) -> Result<T, E>
|
||||
where
|
||||
F: FnMut(u32) -> Fut,
|
||||
Fut: Future<Output = Result<T, E>> + Send,
|
||||
T: Send,
|
||||
E: Display + Send,
|
||||
{
|
||||
let total = policy.attempts;
|
||||
let mut attempt = 1;
|
||||
|
||||
loop {
|
||||
match f(attempt).await {
|
||||
Ok(v) => {
|
||||
if attempt > 1 {
|
||||
logger.log("info", format!("{op} succeeded on attempt {attempt}/{total}"));
|
||||
}
|
||||
return Ok(v);
|
||||
}
|
||||
Err(e) if attempt < total => {
|
||||
let delay = policy.delay(attempt);
|
||||
logger.log(
|
||||
"warn",
|
||||
format!(
|
||||
"{op} attempt {attempt}/{total} failed: {e} - retrying in {}ms",
|
||||
delay.as_millis()
|
||||
),
|
||||
);
|
||||
tokio::time::sleep(delay).await;
|
||||
attempt += 1;
|
||||
}
|
||||
Err(e) => {
|
||||
return Err(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -79,7 +79,12 @@ pub async fn execute_task(
|
||||
let ctx = Arc::new(Context::new());
|
||||
let config_service = ConfigService::new(ctx.clone());
|
||||
let backup_service = BackupService::new(ctx.clone());
|
||||
let config = config_service.load(None).unwrap();
|
||||
|
||||
let local = config_service.load_optional(None);
|
||||
let cache_path = std::path::PathBuf::from(&crate::settings::CONFIG.data_path)
|
||||
.join("dashboard_databases.json");
|
||||
let dashboard = crate::services::dashboard_config::load_cache(&cache_path);
|
||||
let config = crate::services::dashboard_config::merge(&local.databases, &dashboard);
|
||||
|
||||
let metadata_obj = metadata
|
||||
.into_iter()
|
||||
|
||||
Reference in New Issue
Block a user