mirror of
https://github.com/Portabase/agent.git
synced 2026-09-10 01:57:10 +00:00
Merge pull request #74 from Portabase/feat/docker-volume-backup-restore
feat: docker-volume-backup-restore
This commit is contained in:
@@ -36,7 +36,7 @@ jobs:
|
||||
agent-test bash -c "
|
||||
mkdir -p /app/coverage &&
|
||||
rm -rf /app/target/* /app/coverage/* &&
|
||||
cargo test --verbose &&
|
||||
cargo test --verbose -- --test-threads=2 &&
|
||||
sync
|
||||
"
|
||||
|
||||
|
||||
Generated
+1
@@ -3516,6 +3516,7 @@ dependencies = [
|
||||
"azure_core",
|
||||
"azure_storage_blob",
|
||||
"base64 0.22.1",
|
||||
"bollard",
|
||||
"bytes",
|
||||
"chrono",
|
||||
"cron",
|
||||
|
||||
+2
-1
@@ -35,7 +35,7 @@ rand = "0.9.2"
|
||||
bytes = "1.11.0"
|
||||
async-stream = "0.3.6"
|
||||
uuid = { version = "1.20.0", features = ["v4"] }
|
||||
tokio-util = { version = "0.7.18", features = ["compat"] }
|
||||
tokio-util = { version = "0.7.18", features = ["compat", "io"] }
|
||||
tiberius = { version = "0.12", default-features = false, features = ["rustls", "chrono"] }
|
||||
aws-config = "1.8.13"
|
||||
aws-sdk-s3 = { version = "1.122.0", features = ["behavior-version-latest"] }
|
||||
@@ -57,6 +57,7 @@ testcontainers = "0.27.1"
|
||||
testcontainers-modules = { version = "0.15.0", features = ["postgres", "redis", "valkey", "mysql", "mariadb", "mongo"] }
|
||||
postgres = "0.19.12"
|
||||
url = "2.5.8"
|
||||
bollard = "0.20.0"
|
||||
|
||||
[dev-dependencies]
|
||||
tokio = { version = "1", features = ["full"] }
|
||||
|
||||
@@ -124,6 +124,13 @@
|
||||
"port": 1433,
|
||||
"host": "db-mssql",
|
||||
"generated_id": "16706125-ff7e-4c97-8c83-0adeff214682"
|
||||
},
|
||||
{
|
||||
"name": "Test database 14 - Docker Volume",
|
||||
"type": "docker-volume",
|
||||
"volume_name": "databases_sqlite-data",
|
||||
"generated_id": "16706126-ff7e-4c97-8c83-0adeff214690",
|
||||
"container_name": "db-sqlite"
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
+2
-2
@@ -11,7 +11,7 @@ services:
|
||||
- cargo-git:/usr/local/cargo/git
|
||||
- ./databases.json:/config/config.json
|
||||
#- ./databases.toml:/config/config.toml
|
||||
#- /var/run/docker.sock:/var/run/docker.sock
|
||||
- /var/run/docker.sock:/var/run/docker.sock
|
||||
# - cargo-target:/app/target
|
||||
- databases_sqlite-data:/sqlite-data/workspace/data
|
||||
- ./scripts/sqlite/test-db:/sqlite-data-2/workspace/data
|
||||
@@ -19,7 +19,7 @@ services:
|
||||
APP_ENV: development
|
||||
LOG: debug
|
||||
TZ: "Europe/Paris"
|
||||
EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiNWE2YjcxMDgtMGJhYS00Yjg1LTgwMmMtNTNjNjJiMDAzZDgzIiwibWFzdGVyS2V5QjY0IjoiMUh0djdtWCtYVkJxL0IzUEV2WDlZZjlQeUdVZW5oRHlXemo5THRqNW90WT0ifQ=="
|
||||
EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiNDA1MTA4YzQtMDJjYy00NTlhLTkxNjItODExNTc3NjAzMjhjIiwibWFzdGVyS2V5QjY0IjoiMUh0djdtWCtYVkJxL0IzUEV2WDlZZjlQeUdVZW5oRHlXemo5THRqNW90WT0ifQ=="
|
||||
#CHUNK_SIZE_MB: "1"
|
||||
#POOLING: 1
|
||||
#DATABASES_CONFIG_FILE: "config.toml"
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
use crate::domain::docker_volume::docker::{
|
||||
client, create_helper, remove_helper, resolve_helper_image, start_container, stop_container,
|
||||
};
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use bollard::query_parameters::DownloadFromContainerOptions;
|
||||
use futures_util::StreamExt;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
use tokio::fs::File;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, backup_dir: PathBuf, logger: Arc<JobLogger>) -> Result<PathBuf> {
|
||||
tokio::task::spawn_blocking(move || -> Result<PathBuf> {
|
||||
futures::executor::block_on(async move {
|
||||
logger.log("info", format!("Starting docker-volume backup for {}", cfg.name));
|
||||
|
||||
let docker = client()?;
|
||||
let image = resolve_helper_image(&docker).await?;
|
||||
logger.log("debug", format!("Helper image: {image}"));
|
||||
|
||||
if let Some(name) = &cfg.container_name {
|
||||
logger.log("info", format!("Stopping container {name} for consistent backup"));
|
||||
stop_container(&docker, name).await?;
|
||||
}
|
||||
|
||||
let result = async {
|
||||
let helper = create_helper(&docker, &image, &cfg.volume_name, &cfg.generated_id, true, None).await?;
|
||||
|
||||
let file_path = backup_dir.join(format!("{}.tar", cfg.generated_id));
|
||||
let start = Instant::now();
|
||||
|
||||
let dl_opts = DownloadFromContainerOptions { path: "/vol".to_string() };
|
||||
let mut stream = docker.download_from_container(&helper.id, Some(dl_opts));
|
||||
|
||||
let mut out = File::create(&file_path)
|
||||
.await
|
||||
.with_context(|| format!("Failed to create backup file {}", file_path.display()))?;
|
||||
let mut bytes_written: u64 = 0;
|
||||
while let Some(chunk) = stream.next().await {
|
||||
let chunk = chunk.context("Error streaming volume archive from Docker")?;
|
||||
bytes_written += chunk.len() as u64;
|
||||
out.write_all(&chunk).await?;
|
||||
}
|
||||
out.flush().await?;
|
||||
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
logger.log_command("docker download_from_container", None, Some(0), Some(duration_ms));
|
||||
logger.log("info", format!("Volume backup wrote {bytes_written} bytes to {}", file_path.display()));
|
||||
|
||||
remove_helper(&docker, &helper.id).await;
|
||||
anyhow::Ok(file_path)
|
||||
}
|
||||
.await;
|
||||
|
||||
if let Some(name) = &cfg.container_name {
|
||||
if let Err(e) = start_container(&docker, name).await {
|
||||
logger.log("error", format!("Failed to restart container {name}: {e}"));
|
||||
}
|
||||
}
|
||||
|
||||
result
|
||||
})
|
||||
})
|
||||
.await?
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
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};
|
||||
|
||||
pub struct DockerVolumeDatabase {
|
||||
cfg: DatabaseConfig,
|
||||
}
|
||||
|
||||
impl DockerVolumeDatabase {
|
||||
pub fn new(cfg: DatabaseConfig) -> Self {
|
||||
Self { cfg }
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl Database for DockerVolumeDatabase {
|
||||
fn file_extension(&self) -> &'static str {
|
||||
".tar"
|
||||
}
|
||||
|
||||
async fn ping(&self) -> Result<bool> {
|
||||
ping::run(self.cfg.clone()).await
|
||||
}
|
||||
|
||||
async fn backup(&self, dir: &Path, logger: Arc<JobLogger>) -> Result<PathBuf> {
|
||||
FileLock::acquire(&self.cfg.generated_id, DbOpLock::Backup.as_str()).await?;
|
||||
let res = backup::run(self.cfg.clone(), dir.to_path_buf(), logger).await;
|
||||
FileLock::release(&self.cfg.generated_id).await?;
|
||||
res
|
||||
}
|
||||
|
||||
async fn restore(&self, file: &Path, logger: Arc<JobLogger>) -> Result<()> {
|
||||
FileLock::acquire(&self.cfg.generated_id, DbOpLock::Restore.as_str()).await?;
|
||||
let res = restore::run(self.cfg.clone(), file.to_path_buf(), logger).await;
|
||||
FileLock::release(&self.cfg.generated_id).await?;
|
||||
res
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,169 @@
|
||||
#![allow(dead_code)]
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
use bollard::Docker;
|
||||
use bollard::models::{ContainerCreateBody, HostConfig};
|
||||
use bollard::query_parameters::{
|
||||
CreateContainerOptions, InspectContainerOptions, ListContainersOptions,
|
||||
RemoveContainerOptions, StartContainerOptions, StopContainerOptions,
|
||||
};
|
||||
use std::collections::HashMap;
|
||||
use tracing::{info, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
pub const EPHEMERAL_LABEL: &str = "io.portabase.ephemeral";
|
||||
const HELPER_MOUNT: &str = "/vol";
|
||||
|
||||
pub fn client() -> Result<Docker> {
|
||||
Docker::connect_with_unix_defaults().context("Failed to connect to Docker daemon socket")
|
||||
}
|
||||
|
||||
pub fn parse_container_id(mountinfo: &str, cgroup: &str) -> Option<String> {
|
||||
for src in [mountinfo, cgroup] {
|
||||
for line in src.lines() {
|
||||
for marker in ["/containers/", "/docker/"] {
|
||||
if let Some(idx) = line.find(marker) {
|
||||
let rest = &line[idx + marker.len()..];
|
||||
let id: String = rest.chars().take_while(|c| c.is_ascii_hexdigit()).collect();
|
||||
if id.len() >= 64 {
|
||||
return Some(id[..64].to_string());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
pub async fn resolve_helper_image(docker: &Docker) -> Result<String> {
|
||||
if let Ok(img) = std::env::var("PORTABASE_HELPER_IMAGE") {
|
||||
if !img.trim().is_empty() {
|
||||
return Ok(img);
|
||||
}
|
||||
}
|
||||
let mountinfo = std::fs::read_to_string("/proc/self/mountinfo").unwrap_or_default();
|
||||
let cgroup = std::fs::read_to_string("/proc/self/cgroup").unwrap_or_default();
|
||||
let id = parse_container_id(&mountinfo, &cgroup).context(
|
||||
"Could not determine own container id; set PORTABASE_HELPER_IMAGE to a locally-present image",
|
||||
)?;
|
||||
let info = docker
|
||||
.inspect_container(&id, None::<InspectContainerOptions>)
|
||||
.await
|
||||
.with_context(|| format!("Failed to inspect self container {id}"))?;
|
||||
info.image
|
||||
.context("Self container inspection returned no image reference")
|
||||
}
|
||||
|
||||
pub struct Helper {
|
||||
pub id: String,
|
||||
}
|
||||
|
||||
pub async fn create_helper(
|
||||
docker: &Docker,
|
||||
image: &str,
|
||||
volume_name: &str,
|
||||
generated_id: &str,
|
||||
read_only: bool,
|
||||
cmd: Option<Vec<String>>,
|
||||
) -> Result<Helper> {
|
||||
let bind = format!(
|
||||
"{volume_name}:{HELPER_MOUNT}{}",
|
||||
if read_only { ":ro" } else { "" }
|
||||
);
|
||||
let mut labels = HashMap::new();
|
||||
labels.insert(EPHEMERAL_LABEL.to_string(), "true".to_string());
|
||||
labels.insert("com.docker.compose.project".to_string(), String::new());
|
||||
labels.insert("com.docker.compose.service".to_string(), String::new());
|
||||
labels.insert("com.docker.compose.oneoff".to_string(), String::new());
|
||||
|
||||
let name = format!(
|
||||
"portabase-vol-{generated_id}-{}",
|
||||
&Uuid::new_v4().to_string()[..8]
|
||||
);
|
||||
|
||||
let body = ContainerCreateBody {
|
||||
image: Some(image.to_string()),
|
||||
cmd,
|
||||
labels: Some(labels),
|
||||
host_config: Some(HostConfig {
|
||||
binds: Some(vec![bind]),
|
||||
auto_remove: Some(false),
|
||||
..Default::default()
|
||||
}),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let opts = CreateContainerOptions {
|
||||
name: Some(name),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let res = docker
|
||||
.create_container(Some(opts), body)
|
||||
.await
|
||||
.with_context(|| format!("Failed to create helper container for volume {volume_name}"))?;
|
||||
|
||||
Ok(Helper { id: res.id })
|
||||
}
|
||||
|
||||
|
||||
pub async fn remove_helper(docker: &Docker, id: &str) {
|
||||
let stop_opts = StopContainerOptions {
|
||||
t: Some(2),
|
||||
..Default::default()
|
||||
};
|
||||
let _ = docker.stop_container(id, Some(stop_opts)).await;
|
||||
|
||||
if let Ok(info) = docker
|
||||
.inspect_container(id, None::<InspectContainerOptions>)
|
||||
.await
|
||||
{
|
||||
let name = info.name.unwrap_or_default();
|
||||
let name = name.trim_start_matches('/');
|
||||
let code = info.state.and_then(|s| s.exit_code).unwrap_or_default();
|
||||
info!("Helper container {name} exited with code {code}");
|
||||
}
|
||||
|
||||
let opts = RemoveContainerOptions {
|
||||
force: true,
|
||||
..Default::default()
|
||||
};
|
||||
if let Err(e) = docker.remove_container(id, Some(opts)).await {
|
||||
warn!("Failed to remove helper container {id}: {e}");
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn stop_container(docker: &Docker, name: &str) -> Result<()> {
|
||||
docker
|
||||
.stop_container(name, None::<StopContainerOptions>)
|
||||
.await
|
||||
.with_context(|| format!("Failed to stop container {name}"))
|
||||
}
|
||||
|
||||
pub async fn start_container(docker: &Docker, name: &str) -> Result<()> {
|
||||
docker
|
||||
.start_container(name, None::<StartContainerOptions>)
|
||||
.await
|
||||
.with_context(|| format!("Failed to start container {name}"))
|
||||
}
|
||||
|
||||
pub async fn sweep_ephemeral(docker: &Docker) -> Result<usize> {
|
||||
let mut filters = HashMap::new();
|
||||
filters.insert("label".to_string(), vec![format!("{EPHEMERAL_LABEL}=true")]);
|
||||
|
||||
let opts = ListContainersOptions {
|
||||
all: true,
|
||||
filters: Some(filters),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let list = docker.list_containers(Some(opts)).await?;
|
||||
let mut removed = 0;
|
||||
for c in list {
|
||||
if let Some(id) = c.id {
|
||||
remove_helper(docker, &id).await;
|
||||
removed += 1;
|
||||
}
|
||||
}
|
||||
Ok(removed)
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
pub mod backup;
|
||||
pub mod database;
|
||||
pub mod docker;
|
||||
pub mod ping;
|
||||
pub mod restore;
|
||||
@@ -0,0 +1,12 @@
|
||||
use crate::domain::docker_volume::docker::client;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::Result;
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig) -> Result<bool> {
|
||||
let docker = client()?;
|
||||
match docker.inspect_volume(&cfg.volume_name).await {
|
||||
Ok(_) => Ok(true),
|
||||
Err(bollard::errors::Error::DockerResponseServerError { status_code: 404, .. }) => Ok(false),
|
||||
Err(e) => Err(e.into()),
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,102 @@
|
||||
use crate::domain::docker_volume::docker::{
|
||||
client, create_helper, remove_helper, resolve_helper_image, start_container, stop_container,
|
||||
};
|
||||
use crate::services::backup::logger::JobLogger;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use anyhow::{Context, Result};
|
||||
use bollard::exec::StartExecResults;
|
||||
use bollard::models::ExecConfig;
|
||||
use bollard::query_parameters::UploadToContainerOptions;
|
||||
use futures_util::StreamExt;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use std::time::Instant;
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, archive: PathBuf, logger: Arc<JobLogger>) -> Result<()> {
|
||||
tokio::task::spawn_blocking(move || -> Result<()> {
|
||||
futures::executor::block_on(async move {
|
||||
logger.log("info", format!("Starting docker-volume restore for {}", cfg.name));
|
||||
|
||||
let docker = client()?;
|
||||
let image = resolve_helper_image(&docker).await?;
|
||||
|
||||
logger.log("debug", format!("Restore archive: {}", archive.display()));
|
||||
|
||||
if let Some(name) = &cfg.container_name {
|
||||
logger.log("info", format!("Stopping container {name} for restore"));
|
||||
stop_container(&docker, name).await?;
|
||||
}
|
||||
|
||||
let result = async {
|
||||
let helper = create_helper(
|
||||
&docker,
|
||||
&image,
|
||||
&cfg.volume_name,
|
||||
&cfg.generated_id,
|
||||
false,
|
||||
Some(vec![
|
||||
"sh".into(),
|
||||
"-c".into(),
|
||||
"trap 'exit 0' TERM; sleep 2147483647 & wait".into(),
|
||||
]),
|
||||
)
|
||||
.await?;
|
||||
start_container(&docker, &helper.id).await?;
|
||||
|
||||
|
||||
let exec = docker
|
||||
.create_exec(
|
||||
&helper.id,
|
||||
ExecConfig {
|
||||
cmd: Some(vec![
|
||||
"sh".to_string(),
|
||||
"-c".to_string(),
|
||||
"rm -rf /vol/* /vol/.[!.]* 2>/dev/null || true".to_string(),
|
||||
]),
|
||||
attach_stdout: Some(true),
|
||||
attach_stderr: Some(true),
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.context("Failed to create wipe exec")?;
|
||||
|
||||
if let StartExecResults::Attached { mut output, .. } =
|
||||
docker.start_exec(&exec.id, None).await.context("Failed to run wipe exec")?
|
||||
{
|
||||
while output.next().await.is_some() {}
|
||||
}
|
||||
|
||||
let start = Instant::now();
|
||||
|
||||
let file = tokio::fs::File::open(&archive)
|
||||
.await
|
||||
.with_context(|| format!("Failed to open {}", archive.display()))?;
|
||||
let stream = tokio_util::io::ReaderStream::new(file);
|
||||
|
||||
let up_opts = UploadToContainerOptions { path: "/".to_string(), ..Default::default() };
|
||||
docker
|
||||
.upload_to_container(&helper.id, Some(up_opts), bollard::body_try_stream(stream))
|
||||
.await
|
||||
.context("Failed to upload volume archive")?;
|
||||
|
||||
let duration_ms = start.elapsed().as_millis() as f64;
|
||||
logger.log_command("docker upload_to_container", None, Some(0), Some(duration_ms));
|
||||
|
||||
remove_helper(&docker, &helper.id).await;
|
||||
logger.log("info", format!("Volume restore completed for {}", cfg.name));
|
||||
anyhow::Ok(())
|
||||
}
|
||||
.await;
|
||||
|
||||
if let Some(name) = &cfg.container_name {
|
||||
if let Err(e) = start_container(&docker, name).await {
|
||||
logger.log("error", format!("Failed to restart container {name}: {e}"));
|
||||
}
|
||||
}
|
||||
|
||||
result
|
||||
})
|
||||
})
|
||||
.await?
|
||||
}
|
||||
@@ -1,3 +1,4 @@
|
||||
use crate::domain::docker_volume::database::DockerVolumeDatabase;
|
||||
use crate::domain::mongodb::database::MongoDatabase;
|
||||
use crate::domain::mysql::database::MySQLDatabase;
|
||||
use crate::domain::postgres::cluster::database::PostgresClusterDatabase;
|
||||
@@ -41,6 +42,7 @@ impl DatabaseFactory {
|
||||
DbType::Valkey => Arc::new(ValkeyDatabase::new(cfg)),
|
||||
DbType::Firebird => Arc::new(FirebirdDatabase::new(cfg)),
|
||||
DbType::Mssql => Arc::new(MssqlDatabase::new(cfg)),
|
||||
DbType::DockerVolume => Arc::new(DockerVolumeDatabase::new(cfg)),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -59,6 +61,7 @@ impl DatabaseFactory {
|
||||
DbType::Valkey => Arc::new(ValkeyDatabase::new(cfg)),
|
||||
DbType::Firebird => Arc::new(FirebirdDatabase::new(cfg)),
|
||||
DbType::Mssql => Arc::new(MssqlDatabase::new(cfg)),
|
||||
DbType::DockerVolume => Arc::new(DockerVolumeDatabase::new(cfg)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
pub mod docker_volume;
|
||||
pub mod factory;
|
||||
mod mongodb;
|
||||
pub mod mysql;
|
||||
|
||||
+10
@@ -22,6 +22,16 @@ async fn main() {
|
||||
eprintln!("Failed to clean locks on startup: {:?}", e);
|
||||
}
|
||||
|
||||
// Best-effort cleanup of ephemeral helper containers orphaned by a crash.
|
||||
match crate::domain::docker_volume::docker::client() {
|
||||
Ok(docker) => match crate::domain::docker_volume::docker::sweep_ephemeral(&docker).await {
|
||||
Ok(n) if n > 0 => tracing::info!("Removed {n} orphaned ephemeral helper container(s)"),
|
||||
Ok(_) => {}
|
||||
Err(e) => tracing::warn!("Ephemeral helper sweep failed: {e}"),
|
||||
},
|
||||
Err(e) => tracing::debug!("Docker socket unavailable, skipping helper sweep: {e}"),
|
||||
}
|
||||
|
||||
tokio::join!(ping_server(), async {
|
||||
let conn = redis_client::redis_connection().await;
|
||||
scheduler::scheduler_loop(conn).await;
|
||||
|
||||
+20
-3
@@ -26,6 +26,8 @@ pub enum DbType {
|
||||
Valkey,
|
||||
Firebird,
|
||||
Mssql,
|
||||
#[serde(rename = "docker-volume")]
|
||||
DockerVolume,
|
||||
}
|
||||
|
||||
impl DbType {
|
||||
@@ -41,6 +43,7 @@ impl DbType {
|
||||
DbType::Valkey => "valkey",
|
||||
DbType::Firebird => "firebird",
|
||||
DbType::Mssql => "mssql",
|
||||
DbType::DockerVolume => "docker-volume",
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -59,6 +62,8 @@ pub struct DatabaseConfig {
|
||||
pub generated_id: String,
|
||||
pub path: String,
|
||||
pub max_packet_size: String,
|
||||
pub volume_name: String,
|
||||
pub container_name: Option<String>,
|
||||
pub options: HashMap<String, serde_json::Value>,
|
||||
}
|
||||
|
||||
@@ -82,6 +87,8 @@ pub struct InputDatabaseConfig {
|
||||
pub generated_id: String,
|
||||
pub path: Option<String>,
|
||||
pub max_packet_size: Option<String>,
|
||||
pub volume_name: Option<String>,
|
||||
pub container_name: Option<String>,
|
||||
pub options: Option<HashMap<String, serde_json::Value>>,
|
||||
}
|
||||
|
||||
@@ -202,7 +209,7 @@ impl ConfigService {
|
||||
| DbType::Firebird
|
||||
| DbType::Valkey
|
||||
| DbType::Mssql => required(&db.host, &db.name, "host")?,
|
||||
DbType::Sqlite => optional(&db.host),
|
||||
DbType::Sqlite | DbType::DockerVolume => optional(&db.host),
|
||||
};
|
||||
|
||||
let port = match db.db_type {
|
||||
@@ -215,11 +222,13 @@ impl ConfigService {
|
||||
| DbType::Firebird
|
||||
| DbType::Valkey
|
||||
| DbType::Mssql => required(&db.port, &db.name, "port")?,
|
||||
DbType::Sqlite => db.port.unwrap_or(0),
|
||||
DbType::Sqlite | DbType::DockerVolume => db.port.unwrap_or(0),
|
||||
};
|
||||
|
||||
let database_name = match db.db_type {
|
||||
DbType::Sqlite | DbType::Redis | DbType::Valkey => optional(&db.database),
|
||||
DbType::Sqlite | DbType::Redis | DbType::Valkey | DbType::DockerVolume => {
|
||||
optional(&db.database)
|
||||
}
|
||||
DbType::PostgresqlCluster => db
|
||||
.database
|
||||
.clone()
|
||||
@@ -239,6 +248,12 @@ impl ConfigService {
|
||||
_ => String::new(),
|
||||
};
|
||||
|
||||
let volume_name = match db.db_type {
|
||||
DbType::DockerVolume => required(&db.volume_name, &db.name, "volume_name")?,
|
||||
_ => optional(&db.volume_name),
|
||||
};
|
||||
let container_name = db.container_name.clone();
|
||||
|
||||
databases.push(DatabaseConfig {
|
||||
name: db.name,
|
||||
database: database_name,
|
||||
@@ -250,6 +265,8 @@ impl ConfigService {
|
||||
generated_id: db.generated_id,
|
||||
path: path_val,
|
||||
max_packet_size,
|
||||
volume_name,
|
||||
container_name,
|
||||
options: db.options.unwrap_or_default(),
|
||||
});
|
||||
}
|
||||
|
||||
@@ -1,18 +1,19 @@
|
||||
use super::service::RestoreService;
|
||||
|
||||
use crate::utils::compress::decompress_large_tar_gz;
|
||||
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;
|
||||
use crate::services::config::DbType;
|
||||
use crate::utils::common::choose_restore_path;
|
||||
|
||||
impl RestoreService {
|
||||
pub async fn prepare_archive(
|
||||
&self,
|
||||
downloaded_file: PathBuf,
|
||||
tmp_path: &Path,
|
||||
db_type: &DbType,
|
||||
logger: Arc<JobLogger>
|
||||
) -> Result<PathBuf> {
|
||||
logger.log("info", "Start preparing backup archive".to_string());
|
||||
@@ -59,6 +60,13 @@ impl RestoreService {
|
||||
archive = decrypted;
|
||||
}
|
||||
|
||||
if matches!(db_type, DbType::DockerVolume) {
|
||||
let raw_tar = tmp_path.join("volume.tar");
|
||||
crate::utils::compress::gunzip_to_file(archive.as_path(), &raw_tar).await?;
|
||||
logger.log("info", format!("Docker volume archive gunzipped to {}", raw_tar.display()));
|
||||
return Ok(raw_tar);
|
||||
}
|
||||
|
||||
logger.log("info", format!("Decompressing archive {}", archive.display()));
|
||||
|
||||
let files = match decompress_large_tar_gz(archive.as_path(), tmp_path).await {
|
||||
@@ -76,12 +84,8 @@ impl RestoreService {
|
||||
|
||||
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)
|
||||
}
|
||||
let chosen = choose_restore_path(&files, tmp_path, &archive);
|
||||
logger.log("debug", format!("Restore source resolved to: {}", chosen.display()));
|
||||
Ok(chosen)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -27,7 +27,9 @@ impl RestoreService {
|
||||
.download_backup(&file_url, tmp_path, Arc::clone(&logger), expected_size)
|
||||
.await?;
|
||||
|
||||
let backup_file = self.prepare_archive(downloaded, tmp_path, Arc::clone(&logger)).await?;
|
||||
let backup_file = self
|
||||
.prepare_archive(downloaded, tmp_path, &cfg.db_type, Arc::clone(&logger))
|
||||
.await?;
|
||||
|
||||
let result = self.run_restore(cfg, backup_file, Arc::clone(&logger)).await?;
|
||||
|
||||
|
||||
@@ -14,6 +14,8 @@ fn cluster_config() -> DatabaseConfig {
|
||||
generated_id: "40875631-e3d2-4dfe-a26b-2a347ecc64fd".to_string(),
|
||||
path: String::new(),
|
||||
max_packet_size: String::new(),
|
||||
volume_name: String::new(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -36,6 +36,8 @@ async fn start_cluster(user: &str) -> (ContainerAsync<Postgres>, DatabaseConfig)
|
||||
generated_id: "40875631-e3d2-4dfe-a26b-2a347ecc64fd".to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
};
|
||||
(container, config)
|
||||
|
||||
@@ -0,0 +1,279 @@
|
||||
use crate::domain::docker_volume::docker::parse_container_id;
|
||||
|
||||
static ENV_GUARD: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
|
||||
|
||||
#[test]
|
||||
fn parse_container_id_from_mountinfo_line() {
|
||||
let id = "a".repeat(64);
|
||||
let mountinfo = format!(
|
||||
"1234 1000 0:50 / /etc/hostname rw shared:1 - ext4 /var/lib/docker/containers/{id}/hostname rw"
|
||||
);
|
||||
assert_eq!(parse_container_id(&mountinfo, ""), Some(id));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_container_id_from_cgroup_v1() {
|
||||
let id = "b".repeat(64);
|
||||
let cgroup = format!("12:memory:/docker/{id}\n11:cpu:/docker/{id}\n");
|
||||
assert_eq!(parse_container_id("", &cgroup), Some(id));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_container_id_none_on_cgroup_v2() {
|
||||
assert_eq!(parse_container_id("", "0::/\n"), None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn docker_volume_ping_true_for_existing_volume() {
|
||||
use crate::domain::docker_volume::docker::client;
|
||||
use bollard::models::VolumeCreateRequest;
|
||||
use bollard::query_parameters::RemoveVolumeOptions;
|
||||
|
||||
let docker = client().expect("docker daemon required for this test");
|
||||
let vol = format!("portabase-test-{}", uuid::Uuid::new_v4());
|
||||
docker
|
||||
.create_volume(VolumeCreateRequest { name: Some(vol.clone()), ..Default::default() })
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let cfg = volume_config(&vol);
|
||||
let reachable = crate::domain::docker_volume::ping::run(cfg).await.unwrap();
|
||||
assert!(reachable);
|
||||
|
||||
let missing = volume_config("portabase-does-not-exist-xyz");
|
||||
assert!(!crate::domain::docker_volume::ping::run(missing).await.unwrap());
|
||||
|
||||
docker.remove_volume(&vol, None::<RemoveVolumeOptions>).await.ok();
|
||||
}
|
||||
|
||||
fn volume_config(volume_name: &str) -> crate::services::config::DatabaseConfig {
|
||||
use crate::services::config::{DatabaseConfig, DbType};
|
||||
DatabaseConfig {
|
||||
name: "vol-test".to_string(),
|
||||
database: "".to_string(),
|
||||
db_type: DbType::DockerVolume,
|
||||
username: "".to_string(),
|
||||
password: "".to_string(),
|
||||
port: 0,
|
||||
host: "".to_string(),
|
||||
generated_id: uuid::Uuid::new_v4().to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "".to_string(),
|
||||
volume_name: volume_name.to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
async fn ensure_image(docker: &bollard::Docker, image: &str) {
|
||||
use bollard::query_parameters::CreateImageOptionsBuilder;
|
||||
use futures_util::StreamExt;
|
||||
|
||||
let (name, tag) = image.split_once(':').unwrap_or((image, "latest"));
|
||||
let opts = CreateImageOptionsBuilder::default()
|
||||
.from_image(name)
|
||||
.tag(tag)
|
||||
.build();
|
||||
let mut stream = docker.create_image(Some(opts), None, None);
|
||||
while let Some(item) = stream.next().await {
|
||||
item.unwrap();
|
||||
}
|
||||
}
|
||||
|
||||
async fn seed_volume(docker: &bollard::Docker, volume: &str, filename: &str, content: &str) {
|
||||
use bollard::models::{ContainerCreateBody, HostConfig};
|
||||
use bollard::query_parameters::{
|
||||
CreateContainerOptions, RemoveContainerOptions, StartContainerOptions,
|
||||
WaitContainerOptions,
|
||||
};
|
||||
use futures_util::StreamExt;
|
||||
|
||||
ensure_image(docker, "busybox").await;
|
||||
|
||||
let body = ContainerCreateBody {
|
||||
image: Some("busybox".to_string()),
|
||||
cmd: Some(vec![
|
||||
"sh".into(),
|
||||
"-c".into(),
|
||||
format!("printf '%s' '{content}' > /vol/{filename}"),
|
||||
]),
|
||||
host_config: Some(HostConfig {
|
||||
binds: Some(vec![format!("{volume}:/vol")]),
|
||||
..Default::default()
|
||||
}),
|
||||
..Default::default()
|
||||
};
|
||||
let created = docker
|
||||
.create_container(None::<CreateContainerOptions>, body)
|
||||
.await
|
||||
.unwrap();
|
||||
docker.start_container(&created.id, None::<StartContainerOptions>).await.unwrap();
|
||||
let mut wait = docker.wait_container(&created.id, None::<WaitContainerOptions>);
|
||||
while wait.next().await.is_some() {}
|
||||
docker
|
||||
.remove_container(&created.id, Some(RemoveContainerOptions { force: true, ..Default::default() }))
|
||||
.await
|
||||
.ok();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn docker_volume_backup_captures_files() {
|
||||
use crate::domain::docker_volume::docker::client;
|
||||
use bollard::models::VolumeCreateRequest;
|
||||
use bollard::query_parameters::RemoveVolumeOptions;
|
||||
|
||||
let _env_guard = ENV_GUARD.lock().await;
|
||||
unsafe { std::env::set_var("PORTABASE_HELPER_IMAGE", "busybox"); }
|
||||
|
||||
let docker = client().expect("docker daemon required");
|
||||
let vol = format!("portabase-test-{}", uuid::Uuid::new_v4());
|
||||
docker
|
||||
.create_volume(VolumeCreateRequest { name: Some(vol.clone()), ..Default::default() })
|
||||
.await
|
||||
.unwrap();
|
||||
seed_volume(&docker, &vol, "hello.txt", "backup-me").await;
|
||||
|
||||
let tmp = tempfile::TempDir::new().unwrap();
|
||||
let cfg = volume_config(&vol);
|
||||
let logger = std::sync::Arc::new(crate::services::backup::logger::JobLogger::new());
|
||||
let tar = crate::domain::docker_volume::backup::run(cfg, tmp.path().to_path_buf(), logger)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert!(tar.is_file());
|
||||
let names = tar_entry_names(&tar).await;
|
||||
assert!(names.iter().any(|n| n.ends_with("hello.txt")), "entries: {names:?}");
|
||||
|
||||
docker.remove_volume(&vol, None::<RemoveVolumeOptions>).await.ok();
|
||||
}
|
||||
|
||||
async fn tar_entry_names(tar_path: &std::path::Path) -> Vec<String> {
|
||||
use tokio_stream::StreamExt;
|
||||
let f = tokio::fs::File::open(tar_path).await.unwrap();
|
||||
let mut archive = tokio_tar::Archive::new(f);
|
||||
let mut names = Vec::new();
|
||||
let mut entries = archive.entries().unwrap();
|
||||
while let Some(e) = entries.next().await {
|
||||
let e = e.unwrap();
|
||||
names.push(e.path().unwrap().to_string_lossy().to_string());
|
||||
}
|
||||
names
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn docker_volume_restore_is_clean_replace() {
|
||||
use crate::domain::docker_volume::docker::client;
|
||||
use bollard::models::VolumeCreateRequest;
|
||||
use bollard::query_parameters::RemoveVolumeOptions;
|
||||
|
||||
let _env_guard = ENV_GUARD.lock().await;
|
||||
unsafe { std::env::set_var("PORTABASE_HELPER_IMAGE", "busybox"); }
|
||||
|
||||
let docker = client().expect("docker daemon required");
|
||||
let vol = format!("portabase-test-{}", uuid::Uuid::new_v4());
|
||||
docker
|
||||
.create_volume(VolumeCreateRequest { name: Some(vol.clone()), ..Default::default() })
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
seed_volume(&docker, &vol, "keeper.txt", "original").await;
|
||||
|
||||
let tmp = tempfile::TempDir::new().unwrap();
|
||||
let logger = std::sync::Arc::new(crate::services::backup::logger::JobLogger::new());
|
||||
let tar = crate::domain::docker_volume::backup::run(
|
||||
volume_config(&vol),
|
||||
tmp.path().to_path_buf(),
|
||||
logger.clone(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
seed_volume(&docker, &vol, "drift.txt", "added-later").await;
|
||||
|
||||
// Restore uploads the raw Docker tar directly.
|
||||
crate::domain::docker_volume::restore::run(volume_config(&vol), tar.clone(), logger)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let listing = list_volume(&docker, &vol).await;
|
||||
assert!(listing.contains("keeper.txt"), "listing: {listing}");
|
||||
assert!(!listing.contains("drift.txt"), "clean-replace failed, listing: {listing}");
|
||||
|
||||
docker.remove_volume(&vol, None::<RemoveVolumeOptions>).await.ok();
|
||||
}
|
||||
|
||||
async fn list_volume(docker: &bollard::Docker, volume: &str) -> String {
|
||||
use bollard::models::{ContainerCreateBody, HostConfig};
|
||||
use bollard::query_parameters::{
|
||||
CreateContainerOptions, LogsOptions, RemoveContainerOptions, StartContainerOptions,
|
||||
WaitContainerOptions,
|
||||
};
|
||||
use tokio_stream::StreamExt;
|
||||
|
||||
ensure_image(docker, "busybox").await;
|
||||
|
||||
let body = ContainerCreateBody {
|
||||
image: Some("busybox".to_string()),
|
||||
cmd: Some(vec!["sh".into(), "-c".into(), "ls -A /vol".into()]),
|
||||
host_config: Some(HostConfig {
|
||||
binds: Some(vec![format!("{volume}:/vol")]),
|
||||
..Default::default()
|
||||
}),
|
||||
..Default::default()
|
||||
};
|
||||
let created = docker.create_container(None::<CreateContainerOptions>, body).await.unwrap();
|
||||
docker.start_container(&created.id, None::<StartContainerOptions>).await.unwrap();
|
||||
let mut wait = docker.wait_container(&created.id, None::<WaitContainerOptions>);
|
||||
while wait.next().await.is_some() {}
|
||||
|
||||
let mut logs = docker.logs(
|
||||
&created.id,
|
||||
Some(LogsOptions { stdout: true, stderr: false, ..Default::default() }),
|
||||
);
|
||||
let mut out = String::new();
|
||||
while let Some(chunk) = logs.next().await {
|
||||
if let Ok(l) = chunk {
|
||||
out.push_str(&l.to_string());
|
||||
}
|
||||
}
|
||||
docker
|
||||
.remove_container(&created.id, Some(RemoveContainerOptions { force: true, ..Default::default() }))
|
||||
.await
|
||||
.ok();
|
||||
out
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn sweep_removes_labeled_helpers() {
|
||||
use crate::domain::docker_volume::docker::{client, create_helper, sweep_ephemeral, EPHEMERAL_LABEL};
|
||||
use bollard::models::VolumeCreateRequest;
|
||||
use bollard::query_parameters::{ListContainersOptions, RemoveVolumeOptions};
|
||||
use std::collections::HashMap;
|
||||
|
||||
let _env_guard = ENV_GUARD.lock().await;
|
||||
unsafe { std::env::set_var("PORTABASE_HELPER_IMAGE", "busybox"); }
|
||||
let docker = client().expect("docker daemon required");
|
||||
|
||||
let vol = format!("portabase-test-{}", uuid::Uuid::new_v4());
|
||||
docker
|
||||
.create_volume(VolumeCreateRequest { name: Some(vol.clone()), ..Default::default() })
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
ensure_image(&docker, "busybox").await;
|
||||
|
||||
let helper = create_helper(&docker, "busybox", &vol, "sweep-test", true, None).await.unwrap();
|
||||
|
||||
let removed = sweep_ephemeral(&docker).await.unwrap();
|
||||
assert!(removed >= 1);
|
||||
|
||||
let mut filters = HashMap::new();
|
||||
filters.insert("label".to_string(), vec![format!("{EPHEMERAL_LABEL}=true")]);
|
||||
let remaining = docker
|
||||
.list_containers(Some(ListContainersOptions { all: true, filters: Some(filters), ..Default::default() }))
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(remaining.iter().all(|c| c.id.as_deref() != Some(helper.id.as_str())));
|
||||
|
||||
docker.remove_volume(&vol, None::<RemoveVolumeOptions>).await.ok();
|
||||
}
|
||||
@@ -40,6 +40,8 @@ async fn create_config() -> (ContainerAsync<GenericImage>, DatabaseConfig) {
|
||||
generated_id: "3c445eb4-c2c6-4bde-a423-ee1385dcf6d2".to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
};
|
||||
|
||||
|
||||
@@ -32,6 +32,8 @@ async fn create_config() -> (ContainerAsync<Mariadb>, DatabaseConfig) {
|
||||
generated_id: "3c4b4eb4-c2c6-4bde-a423-ee1385dcf6d2".to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "512M".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
};
|
||||
|
||||
|
||||
@@ -7,3 +7,4 @@ mod redis;
|
||||
mod valkey;
|
||||
mod firebird;
|
||||
mod mssql;
|
||||
mod docker_volume;
|
||||
|
||||
@@ -30,6 +30,8 @@ async fn create_config() -> (ContainerAsync<Mongo>, DatabaseConfig) {
|
||||
generated_id: "96d30a9f-ff4b-47c9-aaab-f3147bb34f16".to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
};
|
||||
|
||||
|
||||
@@ -55,6 +55,8 @@ fn make_config(host: String, port: u16, database: &str, generated_id: &str) -> D
|
||||
generated_id: generated_id.to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,6 +32,8 @@ async fn create_config() -> (ContainerAsync<Mysql>, DatabaseConfig) {
|
||||
generated_id: "0f1bb8f2-35a0-4c91-8098-e36873d3ce31".to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "512M".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
};
|
||||
|
||||
|
||||
@@ -39,6 +39,8 @@ async fn create_config() -> (ContainerAsync<Postgres>, DatabaseConfig) {
|
||||
generated_id: "40875631-e3d2-4dfe-a26b-2a347ecc64fd".to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
};
|
||||
|
||||
@@ -157,6 +159,8 @@ async fn postgres_password_with_slash_test() {
|
||||
generated_id: "5a1f0e3c-9b8a-4a8e-9b1b-0a1c2d3e4f5a".to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
};
|
||||
|
||||
|
||||
@@ -29,6 +29,8 @@ async fn create_config() -> (ContainerAsync<Redis>, DatabaseConfig) {
|
||||
generated_id: "40875631-e3d2-4dfe-a26b-2a347ecc64fd".to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
};
|
||||
|
||||
|
||||
@@ -28,6 +28,8 @@ async fn create_config() -> (ContainerAsync<Valkey>, DatabaseConfig) {
|
||||
generated_id: "40875485-e3d2-4dfe-a26b-2a347ecc64fd".to_string(),
|
||||
path: "".to_string(),
|
||||
max_packet_size: "".to_string(),
|
||||
volume_name: "".to_string(),
|
||||
container_name: None,
|
||||
options: std::collections::HashMap::new(),
|
||||
};
|
||||
|
||||
|
||||
@@ -199,3 +199,68 @@ fn keep_ownership_extraction_logic() {
|
||||
let keep4 = opts4.get("keep_ownership").and_then(|v| v.as_bool()).unwrap_or(false);
|
||||
assert!(!keep4, "should strip when value is not bool");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parses_docker_volume_type() {
|
||||
let file = write_json(
|
||||
r#"{
|
||||
"databases": [
|
||||
{
|
||||
"name": "uploads",
|
||||
"type": "docker-volume",
|
||||
"volume_name": "myapp_uploads",
|
||||
"generated_id": "16678159-ff7e-4c97-8c83-0adeff214681",
|
||||
"container_name": "myapp"
|
||||
}
|
||||
]
|
||||
}"#,
|
||||
);
|
||||
|
||||
let service = ConfigService::new(test_context());
|
||||
let cfg = service.load(Some(file.path().to_str().unwrap())).unwrap();
|
||||
|
||||
assert_eq!(cfg.databases[0].db_type.as_str(), "docker-volume");
|
||||
assert_eq!(cfg.databases[0].volume_name, "myapp_uploads");
|
||||
assert_eq!(cfg.databases[0].container_name.as_deref(), Some("myapp"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn docker_volume_container_name_optional() {
|
||||
let file = write_json(
|
||||
r#"{
|
||||
"databases": [
|
||||
{
|
||||
"name": "uploads",
|
||||
"type": "docker-volume",
|
||||
"volume_name": "myapp_uploads",
|
||||
"generated_id": "16678159-ff7e-4c97-8c83-0adeff214681"
|
||||
}
|
||||
]
|
||||
}"#,
|
||||
);
|
||||
|
||||
let service = ConfigService::new(test_context());
|
||||
let cfg = service.load(Some(file.path().to_str().unwrap())).unwrap();
|
||||
|
||||
assert_eq!(cfg.databases[0].volume_name, "myapp_uploads");
|
||||
assert!(cfg.databases[0].container_name.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn docker_volume_requires_volume_name() {
|
||||
let file = write_json(
|
||||
r#"{
|
||||
"databases": [
|
||||
{
|
||||
"name": "uploads",
|
||||
"type": "docker-volume",
|
||||
"generated_id": "16678159-ff7e-4c97-8c83-0adeff214681"
|
||||
}
|
||||
]
|
||||
}"#,
|
||||
);
|
||||
|
||||
let service = ConfigService::new(test_context());
|
||||
let err = service.load(Some(file.path().to_str().unwrap())).unwrap_err();
|
||||
assert!(err.contains("volume_name"), "error was: {err}");
|
||||
}
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
use crate::utils::common::{BackupMethod, vec_to_option_json};
|
||||
use crate::utils::common::{BackupMethod, choose_restore_path, vec_to_option_json};
|
||||
use serde_json::json;
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
#[test]
|
||||
fn backup_method_to_string_automatic() {
|
||||
@@ -47,3 +48,28 @@ fn vec_to_option_json_serializes_struct_vector() {
|
||||
]))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn choose_restore_path_single_file_returns_that_file() {
|
||||
let dir = Path::new("/tmp/extract");
|
||||
let archive = Path::new("/tmp/backup.tar.gz");
|
||||
let files = vec![PathBuf::from("/tmp/extract/dump.sql")];
|
||||
// Single extracted file: restore from that file directly.
|
||||
assert_eq!(
|
||||
choose_restore_path(&files, dir, archive),
|
||||
PathBuf::from("/tmp/extract/dump.sql")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn choose_restore_path_multi_file_non_docker_volume_returns_archive_path() {
|
||||
let dir = Path::new("/tmp/extract");
|
||||
let archive = Path::new("/tmp/backup.tar.gz");
|
||||
let files = vec![
|
||||
PathBuf::from("/tmp/extract/toc.dat"),
|
||||
PathBuf::from("/tmp/extract/3141.dat.gz"),
|
||||
];
|
||||
let chosen = choose_restore_path(&files, dir, archive);
|
||||
assert_eq!(chosen, PathBuf::from("/tmp/backup.tar.gz"));
|
||||
assert_eq!(chosen.extension().and_then(|e| e.to_str()), Some("gz"));
|
||||
}
|
||||
|
||||
@@ -71,3 +71,57 @@ async fn decompress_multiple_files() -> Result<()> {
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn compress_tar_is_not_double_wrapped() -> Result<()> {
|
||||
use tokio_tar::Builder as TarBuilder;
|
||||
|
||||
let tmp = tempdir()?;
|
||||
// Build a real tar containing a single entry "payload.txt".
|
||||
let payload = tmp.path().join("payload.txt");
|
||||
write(&payload, b"volume-bytes").await?;
|
||||
let tar_path = tmp.path().join("volume.tar");
|
||||
{
|
||||
let f = tokio::fs::File::create(&tar_path).await?;
|
||||
let mut b = TarBuilder::new(f);
|
||||
b.append_path_with_name(&payload, "payload.txt").await?;
|
||||
b.finish().await?;
|
||||
}
|
||||
|
||||
let result = compress_to_tar_gz_large(
|
||||
&tar_path,
|
||||
std::sync::Arc::new(crate::services::backup::logger::JobLogger::new()),
|
||||
)
|
||||
.await?;
|
||||
assert_eq!(result.compressed_path, tmp.path().join("volume.tar.gz"));
|
||||
|
||||
// Decompress and confirm the FIRST tar entry is "payload.txt" — i.e. our tar
|
||||
// was gzipped directly, not wrapped inside another tar named "volume.tar".
|
||||
let out = tmp.path().join("out");
|
||||
tokio::fs::create_dir_all(&out).await?;
|
||||
let files = decompress_large_tar_gz(&result.compressed_path, &out).await?;
|
||||
assert_eq!(files.len(), 1);
|
||||
assert_eq!(files[0].file_name().unwrap(), "payload.txt");
|
||||
assert_eq!(read(&files[0]).await?, b"volume-bytes");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn gunzip_to_file_restores_tar_byte_for_byte() -> Result<()> {
|
||||
use crate::utils::compress::gunzip_to_file;
|
||||
|
||||
let tmp = tempdir()?;
|
||||
let tar_path = tmp.path().join("input.tar");
|
||||
let original: Vec<u8> = (0u32..50_000).map(|n| (n % 256) as u8).collect();
|
||||
write(&tar_path, &original).await?;
|
||||
|
||||
let gz = compress_to_tar_gz_large(&tar_path, std::sync::Arc::new(crate::services::backup::logger::JobLogger::new())).await?;
|
||||
|
||||
let out_tar = tmp.path().join("out.tar");
|
||||
gunzip_to_file(&gz.compressed_path, &out_tar).await?;
|
||||
|
||||
assert_eq!(read(&out_tar).await?, original);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -1,5 +1,17 @@
|
||||
use serde::Serialize;
|
||||
use serde_json::Value;
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
pub(crate) fn choose_restore_path(
|
||||
extracted: &[PathBuf],
|
||||
_extraction_dir: &Path,
|
||||
archive: &Path,
|
||||
) -> PathBuf {
|
||||
match extracted.len() {
|
||||
1 => extracted[0].clone(),
|
||||
_ => archive.to_path_buf(),
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
pub enum BackupMethod {
|
||||
|
||||
+54
-1
@@ -1,4 +1,4 @@
|
||||
use anyhow::Result;
|
||||
use anyhow::{Context, Result};
|
||||
use async_compression::tokio::bufread::GzipDecoder;
|
||||
use async_compression::tokio::write::GzipEncoder as AsyncGzipEncoder;
|
||||
use futures::StreamExt;
|
||||
@@ -32,6 +32,39 @@ pub async fn compress_to_tar_gz_large(file: &PathBuf, logger: Arc<JobLogger>) ->
|
||||
});
|
||||
}
|
||||
|
||||
if file
|
||||
.file_name()
|
||||
.and_then(|n| n.to_str())
|
||||
.map(|n| n.ends_with(".tar"))
|
||||
.unwrap_or(false)
|
||||
{
|
||||
let gz_path = PathBuf::from(format!("{}.gz", file.display()));
|
||||
logger.log("info", format!("Input {:?} is a raw tar, gzipping directly", file));
|
||||
|
||||
let input = File::open(file)
|
||||
.await
|
||||
.map_err(|e| anyhow::anyhow!("Failed to open tar {:?}: {}", file, e))?;
|
||||
let mut reader = BufReader::with_capacity(8 * 1024 * 1024, input);
|
||||
|
||||
let output_file = File::create(&gz_path)
|
||||
.await
|
||||
.map_err(|e| anyhow::anyhow!("Failed to create {:?}: {}", gz_path, e))?;
|
||||
let mut encoder = AsyncGzipEncoder::new(output_file);
|
||||
|
||||
tokio::io::copy(&mut reader, &mut encoder)
|
||||
.await
|
||||
.map_err(|e| anyhow::anyhow!("Gzip copy failed: {}", e))?;
|
||||
encoder
|
||||
.shutdown()
|
||||
.await
|
||||
.map_err(|e| anyhow::anyhow!("Gzip shutdown failed: {}", e))?;
|
||||
|
||||
logger.log("info", format!("Compressed {:?} to {:?}", file, gz_path));
|
||||
return Ok(CompressionResult {
|
||||
compressed_path: gz_path,
|
||||
});
|
||||
}
|
||||
|
||||
let tar_gz_path = file.with_extension("").with_extension("tar.gz");
|
||||
|
||||
let output_file = File::create(&tar_gz_path)
|
||||
@@ -119,3 +152,23 @@ pub async fn decompress_large_tar_gz(
|
||||
|
||||
Ok(extracted_files)
|
||||
}
|
||||
|
||||
pub async fn gunzip_to_file(gz_path: &Path, out_path: &Path) -> Result<()> {
|
||||
let file = File::open(gz_path)
|
||||
.await
|
||||
.with_context(|| format!("Failed to open {}", gz_path.display()))?;
|
||||
let buf_reader = BufReader::with_capacity(8 * 1024 * 1024, file);
|
||||
let mut decoder = GzipDecoder::new(buf_reader);
|
||||
|
||||
let out = File::create(out_path)
|
||||
.await
|
||||
.with_context(|| format!("Failed to create {}", out_path.display()))?;
|
||||
let mut writer = tokio::io::BufWriter::new(out);
|
||||
|
||||
tokio::io::copy(&mut decoder, &mut writer).await?;
|
||||
writer.shutdown().await?;
|
||||
|
||||
info!("Gunzipped {:?} into {:?}", gz_path, out_path);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user