Compare commits

..

75 Commits

Author SHA1 Message Date
github-actions[bot] d42c31f3dc chore: release 1.21.1 2026-09-17 11:43:39 +00:00
Charles GTE 84673373a7 Merge pull request #107 from Portabase/fix/mongodb-options
fix: mongodb-options
2026-09-17 13:41:27 +02:00
Charles GTE 6f8a5c4228 fix: mongodb options 2026-09-17 11:20:16 +02:00
charles-gauthereau c5806a6ea4 fix: refactoring 2026-09-17 10:47:32 +02:00
Charles GTE c241020542 fix: mongodb options 2026-09-17 10:13:58 +02:00
github-actions[bot] e0bc6556aa chore: release 1.21.0 2026-09-14 13:11:26 +00:00
Charles GTE 947cb9bed7 Merge pull request #104 from Portabase/dev
Dev
2026-09-14 15:09:04 +02:00
Charles GTE b5d0035a3b Merge pull request #106 from Portabase/poc/hmac
fix: openssl windows
2026-09-14 14:52:26 +02:00
Charles GTE d52c328a99 fix: refactoring 2026-09-14 14:51:58 +02:00
Charles GTE 9b79ac4850 Merge branch 'refs/heads/dev' into poc/hmac
# Conflicts:
#	Cargo.lock
2026-09-14 14:51:22 +02:00
Charles GTE 6e50718b65 fix: tests 2026-09-14 14:30:11 +02:00
Charles GTE ee5ed81d48 fix: tests 2026-09-14 14:12:29 +02:00
Charles GTE df078b6202 Merge pull request #105 from Portabase/feat/sftp-storage
feat: sftp-storage
2026-09-14 10:55:40 +02:00
Charles GTE 24e3e3098d chore: update packages 2026-09-13 17:25:13 +02:00
Charles GTE f6e1a0df98 fix: refactoring 2026-09-12 21:15:04 +02:00
Charles GTE 49366ee24e refactor(storage): generic rclone config builder and shared obscure util 2026-09-12 21:10:50 +02:00
Charles GTE 012a984973 fix(storage): reject line breaks in sftp host/username; drop unused import 2026-09-12 21:02:28 +02:00
Charles GTE 6d1d7c6891 fix: refactoring 2026-09-12 18:15:45 +02:00
Charles GTE c70cde8a9c feat(storage): sftp provider over rclone rcat 2026-09-12 17:57:04 +02:00
Charles GTE c353443bdd feat(storage): sftp config model and rclone-config builder 2026-09-12 17:57:04 +02:00
Charles GTE 9f9dedfe87 feat: rclone-privider (#103)
* build: install rclone 1.75.1 in agent images

* build: exclude target/ and local artifacts from the docker context

* feat(storage): add rclone config model and validation helpers

* feat(storage): stream backups into rclone rcat

* fix(storage): abort truncated rclone uploads and strengthen stderr test

* feat(storage): add RcloneProvider upload path

* fix(storage): bound rclone timeouts, detect restore truncation, track sdd workspace

* fix(storage): block virtual and non-viable rclone backends

* fix: refactoring
2026-09-12 17:09:57 +02:00
Charles GTE 9bb830327e fix: refactoring 2026-09-10 22:55:47 +02:00
Charles GTE 31e55a7b4f fix(storage): block virtual and non-viable rclone backends 2026-09-10 22:38:08 +02:00
Charles GTE 576c055092 fix(storage): bound rclone timeouts, detect restore truncation, track sdd workspace 2026-09-10 22:38:08 +02:00
Charles GTE ba83b6b2b7 feat(storage): add RcloneProvider upload path 2026-09-10 22:38:07 +02:00
Charles GTE ea356ed54f fix(storage): abort truncated rclone uploads and strengthen stderr test 2026-09-10 22:37:37 +02:00
Charles GTE b222b13934 feat(storage): stream backups into rclone rcat 2026-09-10 22:37:09 +02:00
Charles GTE 9409f0ff25 feat(storage): add rclone config model and validation helpers 2026-09-10 22:37:00 +02:00
Charles GTE 59312215e9 build: exclude target/ and local artifacts from the docker context 2026-09-10 22:36:46 +02:00
Charles GTE 3ff360a971 build: install rclone 1.75.1 in agent images 2026-09-10 22:36:38 +02:00
Théo LAGACHE cc1703592e poc 2026-09-09 11:57:56 +02:00
github-actions[bot] 5f92297108 chore: release 1.20.1 2026-09-03 15:41:20 +00:00
Charles GTE d1821f7045 Merge pull request #102 from Portabase/fix/cron-crash
fix: cron crash for databases setup from dashboard
2026-09-03 17:38:53 +02:00
charles-gauthereau 0fe4d50cbd fix: docker-compose.yml 2026-09-03 17:38:35 +02:00
charles-gauthereau 1f29466a28 fix: cron crash for databases setup from dashboard 2026-09-03 17:22:03 +02:00
github-actions[bot] b74aaa0bbc chore: release 1.20.0 2026-08-28 15:42:27 +00:00
Charles GTE 7b5e5b2c78 Merge pull request #101 from Portabase/feat/retry-system
feat: retry-system
2026-08-28 17:40:20 +02:00
charles-gauthereau 99f1ef2081 fix: refactoring 2026-08-28 17:25:30 +02:00
charles-gauthereau 42c3c5945a fix: refactoring 2026-08-28 17:06:23 +02:00
charles-gauthereau db0f87c2e9 Merge branch 'main' into feat/retry-system 2026-08-28 16:54:51 +02:00
github-actions[bot] f95f0aa73a chore: release 1.19.2 2026-08-28 07:59:57 +00:00
Charles GTE 2177dfd44b Merge pull request #100 from Portabase/fix/kill-on-drop
fix: kill_on_drop on firebird, mariadb, mysql, redis, valkey
2026-08-28 09:57:51 +02:00
charles-gauthereau 0a53eec184 fix: kill_on_drop on firebird, mariadb, mysql, redis, valkey 2026-08-28 09:43:08 +02:00
charles-gauthereau 87f33af772 fix: wire retry env vars into helm configmap, dedupe terminal retry log
helm/templates/env-configmap.yaml never listed RETRY_ATTEMPTS and
RETRY_BACKOFF_MS even though values.yaml gained them, so --set
env.RETRY_ATTEMPTS=N was silently ignored by Kubernetes deployments.
Add both keys in the same explicit style as the existing entries.

src/utils/retry.rs logged its own "failed after N attempts" error on
exhaustion, on top of the terminal log each call site already writes,
producing two error entries per failure. Worse, it changed a log
level: FileLock::acquire's "backup_already_in_progress" bails through
the combinator, which now logged it as error before runner.rs got a
chance to reclassify it as the routine warn it always was. A manual
backup colliding with a scheduled one would show up as a hard error
on the dashboard instead of the harmless warn it used to be, breaking
the "fails exactly as it does today" guarantee for job records.

Drop the combinator's terminal error log and give download_backup its
own terminal error log so all three call sites (runner, uploader,
downloader) own their failure logging uniformly. Update the two tests
that asserted the removed message to assert the new behavior instead.
2026-08-27 22:46:49 +02:00
charles-gauthereau b6a120fcbf feat: retry the restore backup download
Extracts the download body to download_once and makes download_backup a
retry wrapper around it. This path had no retry at all before, so a
single dropped connection failed the whole restore job.

Retrying is safe because File::create truncates and the target filename
is derived from Content-Disposition or the URL, so it is stable across
attempts. There is no Range resume: a download that fails at 90% starts
over.
2026-08-27 19:01:10 +02:00
charles-gauthereau 5298d82576 feat: retry storage uploads
Wraps provider.upload in the retry combinator. No provider changes are
needed: each one builds its upload stream from the file on disk inside
upload(), so every attempt gets a fresh handle and a fresh nonce.

Result<UploadResult, UploadResult> is collapsed with an or-pattern so
the last attempt's error and metadata survive into the existing failure
branch. A missing backup file short-circuits into Err rather than
returning early, which skips the retry without skipping the
backup_upload_status(failed) call that closes the server-side record.

backup_upload_init and backup_upload_status are left unwrapped; they
are control-plane calls, not storage uploads.
2026-08-27 18:53:06 +02:00
charles-gauthereau b0da2e40a3 feat: retry the database backup
Each attempt now dumps into its own tmp_path/attempt-{n} directory,
which is removed when the attempt fails. Without the per-attempt
directory pg_dump -Fd would refuse every retry, because it will not
write into a directory a previous attempt left behind; removing it on
failure keeps peak disk at one attempt's artifacts rather than five.

A backup blocked by a concurrent job is retried before surfacing the
same backup_already_in_progress code, since FileLock reports it as an
ordinary error and the combinator cannot tell it apart.

Reshapes retry()'s bound from the native AsyncFnMut sugar to the
classic F: FnMut(u32) -> Fut, Fut: Future<Output = Result<T, E>> + Send
shape (the pattern tokio-retry and backoff both use). AsyncFnMut's
produced future is a lifetime-quantified associated type
(F::CallRefFuture<'_>) that cannot be named or bounded as Send on
stable Rust, so wiring a retried call through it into a future that
eventually gets polled inside tokio::spawn (dispatcher.rs, via
execute_backup) made rustc's opaque-type Send inference fail with
"implementation of Send is not general enough" at the spawn site,
several call layers away from the actual retry call. Naming Fut as its
own type parameter lets Send be asserted on it directly instead, which
resolves cleanly. The combinator's control flow, log messages, and
formats are unchanged; only the bound and the call sites' closure
shape (async |x| { } becomes |x| async { }, with shared references
bound outside a move closure so the inner async move block only moves
Copy references, not the originals) are affected.
2026-08-27 18:39:41 +02:00
charles-gauthereau 55e20d48e7 feat: configurable retry policy and combinator
Adds RETRY_ATTEMPTS (3..=5, default 3) and RETRY_BACKOFF_MS
(100..=30000, default 1000) to Settings, validated with the same
panic-on-invalid contract as POOLING and CHUNK_SIZE_MB.

The combinator logs every failed attempt and any late success through
the JobLogger it borrows, so retries reach the server on the existing
job-log path with no API change. It borrows rather than clones the Arc
so Arc::try_unwrap in the backup executor keeps working. Backoff is
exponential with equal jitter, because the uploader retries storages
concurrently and would otherwise retry them in lockstep.
2026-08-27 18:13:17 +02:00
github-actions[bot] f7f5f7e141 chore: release 1.19.1 2026-08-21 17:09:15 +00:00
Charles GTE 2170f96a72 Merge pull request #99 from Portabase/fix/config-file
fix: config directory creation
2026-08-21 19:06:48 +02:00
Charles GTE 298d46ba81 fix: config directory creation 2026-08-21 13:50:40 +02:00
github-actions[bot] 6537e9df53 chore: release 1.19.0 2026-08-20 11:56:59 +00:00
Charles GTE b8d869d5a6 Merge pull request #97 from Portabase/feat/dashboard-databases
feat: dashboard-databases
2026-08-20 13:54:46 +02:00
Charles GTE 90941ea67d fix 2026-08-20 13:41:14 +02:00
charles-gauthereau 046d593e2b fix 2026-08-18 18:20:21 +02:00
charles-gauthereau cf9a59a138 fix: config.rs 2026-08-18 13:07:22 +02:00
charles-gauthereau 2464f6dfb0 Merge branch 'main' into feat/dashboard-databases
# Conflicts:
#	docker-compose.yml
#	src/services/config.rs
2026-08-18 13:05:56 +02:00
github-actions[bot] 069067ca55 chore: release 1.18.6 2026-08-17 09:10:08 +00:00
Charles GTE 04b654d219 Merge pull request #96 from Portabase/fix/mongo-cloud
fix: mongodb cloud cluster issue
2026-08-17 11:07:35 +02:00
charles-gauthereau d83511cb64 fix: mongodb cloud cluster issue 2026-08-17 10:52:21 +02:00
charles-gauthereau ca294e968c fix 2026-08-13 22:02:17 +02:00
charles-gauthereau 1be88ffdb9 Merge branch 'main' into feat/dashboard-databases
# Conflicts:
#	docker-compose.yml
2026-08-13 21:56:38 +02:00
charles-gauthereau 9d393a96a4 fix 2026-08-13 21:55:59 +02:00
charles-gauthereau 044bf80633 feat(agent): ingest dashboard databases via merged cycle + cache 2026-08-13 18:47:18 +02:00
charles-gauthereau 0a6eb6db22 feat(dashboard_config): atomic cache load/persist 2026-08-13 18:40:56 +02:00
charles-gauthereau b26ff81889 feat(dashboard_config): merge (dashboard-wins) + collect_configs 2026-08-13 18:39:54 +02:00
charles-gauthereau 4e2f29f4ca feat(status): decrypt dashboard config_ciphertext into resolved_config 2026-08-13 18:38:15 +02:00
charles-gauthereau d1c8df4cac refactor(config): extract build_config, add load_optional + Serialize 2026-08-13 18:35:52 +02:00
github-actions[bot] 30f83bafcf chore: release 1.18.5 2026-07-26 09:06:42 +00:00
Charles GTE de106c835e Merge pull request #93 from Portabase/fix/s3-storage
fix: s3 logs
2026-07-26 11:04:35 +02:00
Charles GTE 23d6822ddc fix: s3 logs 2026-07-26 11:04:07 +02:00
github-actions[bot] fe1d74945f chore: release 1.18.4 2026-07-25 14:03:12 +00:00
Charles GTE 424a646385 Merge pull request #92 from Portabase/fix/windows-build
fix: windows-build
2026-07-25 16:00:57 +02:00
Charles GTE 0ef4bba5d7 fix: docker-compose.yml 2026-07-25 15:45:56 +02:00
Charles GTE 6548140eaf fix: windows build 2026-07-25 15:19:46 +02:00
62 changed files with 3692 additions and 1667 deletions
+3 -5
View File
@@ -1,11 +1,9 @@
# Git
.git
.gitignore
# MD files
CHANGELOG.md
README.md
RELEASE.md
#IDE configurations
.idea
target
dump.rdb
.superpowers
+12
View File
@@ -118,12 +118,24 @@ jobs:
secrets:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
build-windows:
needs: create-release
if: ${{ needs.create-release.result == 'success' }}
uses: ./.github/workflows/windows-release.yml
with:
version: ${{ needs.create-release.outputs.version }}
ref: ${{ needs.create-release.outputs.version }}
draft_tag: ${{ needs.create-release.outputs.draft_tag }}
secrets:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
finalize-release:
needs:
- create-release
- publish-docker
- publish-docker-ghcr
- publish-helm
- build-windows
runs-on: ubuntu-latest
outputs:
release_tag: ${{ steps.publish_release_step.outputs.release_tag }}
+41 -52
View File
@@ -1,14 +1,30 @@
name: Build Windows release
on:
workflow_call:
inputs:
version:
description: 'Release version (git tag), e.g. 1.18.4'
type: string
required: false
ref:
description: 'Git ref to check out and build'
type: string
required: false
draft_tag:
description: 'Draft GitHub release tag to attach the asset to (e.g. untagged-xxxx). Empty = skip upload.'
type: string
required: false
secrets:
GH_TOKEN:
required: false
workflow_dispatch:
push:
tags:
- '[0-9]+.[0-9]+.[0-9]+'
branches:
- main
- master
inputs:
ref:
description: 'Git ref to check out and build'
type: string
required: false
jobs:
build-windows:
@@ -17,80 +33,53 @@ jobs:
steps:
- name: Checkout repository
uses: actions/checkout@v4
with:
ref: ${{ inputs.ref || github.ref }}
- name: Set up Rust toolchain (MSVC)
uses: actions-rs/toolchain@v1
uses: dtolnay/rust-toolchain@stable
with:
toolchain: stable-x86_64-pc-windows-msvc
profile: minimal
override: true
targets: x86_64-pc-windows-msvc
- name: Install vcpkg and OpenSSL (x64)
shell: pwsh
run: |
# Install vcpkg and the prebuilt OpenSSL package
git clone https://github.com/microsoft/vcpkg C:\vcpkg
C:\vcpkg\bootstrap-vcpkg.bat
C:\vcpkg\vcpkg install openssl:x64-windows
# Export variables for subsequent steps
'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: Cache cargo build
uses: Swatinem/rust-cache@v2
- name: Build (cargo release)
shell: pwsh
env:
# Cargo / openssl-sys will pick up OPENSSL_DIR from the environment
OPENSSL_DIR: ${{ env.OPENSSL_DIR }}
run: |
# Ensure the environment variable is present for this step
if (-Not $env:OPENSSL_DIR) { Write-Host "OPENSSL_DIR not set, printing env for debugging"; Get-ChildItem Env: | ForEach-Object { Write-Host $_ } }
# Build the declared bin target explicitly (Cargo.toml [[bin]] name = "app")
cargo build --release --bin app
run: cargo build --release --bin app
- name: Prepare artifact zip
id: prepare_artifact
shell: pwsh
env:
RELEASE_TAG: ${{ github.ref_name }}
RELEASE_VERSION: ${{ inputs.version }}
run: |
$tag = $env:RELEASE_TAG
$tag = $env:RELEASE_VERSION
if (-not $tag) { $tag = $env:GITHUB_SHA }
# Package the declared bin target deterministically (Cargo.toml [[bin]] name = "app")
$exe = "target\release\app.exe"
if (-not (Test-Path $exe)) { Write-Error "Built binary $exe not found in target/release"; exit 1 }
$outDir = "artifact"
New-Item -ItemType Directory -Path $outDir -Force | Out-Null
# Ship under the package name, not the internal bin name "app"
# Ship under the package name, not the internal bin name "app".
Copy-Item -Path $exe -Destination "$outDir\portabase-agent.exe"
$zipName = "windows-release-$tag.zip"
if (Test-Path $zipName) { Remove-Item $zipName }
Compress-Archive -Path "$outDir\*" -DestinationPath $zipName -Force
Write-Host "ZIP=$zipName"
Write-Output "zip=$zipName" | Out-File -FilePath $env:GITHUB_OUTPUT -Encoding utf8 -Append
"zip=$zipName" | Out-File -FilePath $env:GITHUB_OUTPUT -Encoding utf8 -Append
- name: Upload build artifact
uses: actions/upload-artifact@v4
with:
name: windows-release
path: windows-release-*.zip
path: ${{ steps.prepare_artifact.outputs.zip }}
- name: Create GitHub Release
if: startsWith(github.ref, 'refs/tags/')
id: create_release
uses: softprops/action-gh-release@v1
with:
tag_name: ${{ github.ref_name }}
- name: Attach asset to draft release
if: ${{ inputs.draft_tag != '' }}
shell: pwsh
env:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
- name: Upload release asset
if: startsWith(github.ref, 'refs/tags/')
uses: actions/upload-release-asset@v1
with:
upload_url: ${{ steps.create_release.outputs.upload_url }}
asset_path: windows-release-${{ github.ref_name }}.zip
asset_name: windows-release-${{ github.ref_name }}.zip
asset_content_type: application/zip
GH_TOKEN: ${{ secrets.GH_TOKEN }}
run: |
gh release upload "${{ inputs.draft_tag }}" "${{ steps.prepare_artifact.outputs.zip }}" --clobber
+2 -1
View File
@@ -6,5 +6,6 @@
.env
.claude
/docs
.superpowers
+1 -1
View File
@@ -27,5 +27,5 @@ keywords:
- self-hosted
- portabase
license: Apache-2.0
version: 1.18.3
version: 1.21.1
date-released: '2026-02-24'
Generated
+1216 -1354
View File
File diff suppressed because it is too large Load Diff
+5 -3
View File
@@ -1,6 +1,6 @@
[package]
name = "portabase-agent"
version = "1.18.3"
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"
@@ -57,6 +58,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"
percent-encoding = "2.3.2"
bollard = "0.20.0"
[dev-dependencies]
+3 -3
View File
@@ -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:
+6 -47
View File
@@ -1,13 +1,14 @@
services:
rust-app:
# build:
# context: .
# dockerfile: docker/Dockerfile
# target: prod
image: portabase/agent:latest
build:
context: .
dockerfile: docker/Dockerfile
target: prod
# image: portabase/agent:latest
container_name: rust-prod
volumes:
- ./databases.json:/config/config.json
- /var/run/docker.sock:/var/run/docker.sock
environment:
LOG: info
TZ: "Europe/Paris"
@@ -18,48 +19,6 @@ services:
networks:
- portabase
db-mongodb-auth:
container_name: db-mongodb-auth
image: mongo:latest
ports:
- "27082:27017"
environment:
MONGO_INITDB_ROOT_USERNAME: root
MONGO_INITDB_ROOT_PASSWORD: rootpassword
MONGO_INITDB_DATABASE: testdbauth
command: mongod --auth
networks:
- portabase
volumes:
- mongodb-data-auth:/data/db
healthcheck:
test: [ "CMD", "mongo", "--eval", "db.adminCommand('ping')" ]
interval: 5s
timeout: 5s
retries: 10
db-mongodb:
container_name: db-mongodb
image: mongo:latest
ports:
- "27083:27017"
volumes:
- mongodb-data:/data/db
healthcheck:
test: [ "CMD", "mongosh", "--eval", "db.adminCommand('ping')" ]
interval: 5s
timeout: 5s
retries: 10
environment:
MONGO_INITDB_DATABASE: testdb
networks:
- portabase
volumes:
mongodb-data:
mongodb-data-auth:
networks:
portabase:
name: portabase_network
+3 -1
View File
@@ -21,9 +21,11 @@ services:
LOG: debug
TZ: "Europe/Paris"
# TMPDIR: /scratch
EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiNmM4NWE3ODQtODRkMi00YzUyLTgzYmUtZTc2MDZkZjg2YjM5IiwibWFzdGVyS2V5QjY0IjoiMUh0djdtWCtYVkJxL0IzUEV2WDlZZjlQeUdVZW5oRHlXemo5THRqNW90WT0ifQ=="
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
View File
@@ -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
+3 -1
View File
@@ -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 }}
+2
View File
@@ -11,6 +11,8 @@ env:
POLLING: "5"
APP_ENV: "production"
LOG: "info"
RETRY_ATTEMPTS: "3"
RETRY_BACKOFF_MS: "1000"
resources:
limits:
+30 -8
View File
@@ -2,13 +2,16 @@
use crate::core::context::Context;
use crate::services::backup::BackupService;
use crate::services::config::ConfigService;
use crate::services::config::{ConfigService, DatabaseConfig};
use crate::services::cron::CronService;
use crate::services::dashboard_config::{collect_configs, load_cache, merge, persist_cache};
use crate::services::restore::RestoreService;
use crate::services::status::StatusService;
use crate::settings::CONFIG;
use crate::utils::common::BackupMethod;
use std::path::PathBuf;
use std::sync::Arc;
use tracing::info;
use tracing::{error, info, warn};
pub struct Agent {
ctx: Arc<Context>,
@@ -17,6 +20,8 @@ pub struct Agent {
cron_service: CronService,
backup_service: BackupService,
restore_service: RestoreService,
dashboard_cache: Vec<DatabaseConfig>,
cache_path: PathBuf,
}
impl Agent {
@@ -28,6 +33,9 @@ impl Agent {
let backup_service = BackupService::new(ctx.clone());
let restore_service = RestoreService::new(ctx.clone());
let cache_path = PathBuf::from(&CONFIG.data_path).join("dashboard_databases.json");
let dashboard_cache = load_cache(&cache_path);
Agent {
ctx,
config_service,
@@ -35,19 +43,33 @@ impl Agent {
cron_service,
backup_service,
restore_service,
dashboard_cache,
cache_path,
}
}
pub async fn run(&mut self, method: BackupMethod) -> Result<(), Box<dyn std::error::Error>> {
let config = self.config_service.load(None)?;
let ping_result = self.status_service.ping(&config.databases).await?;
let local = self.config_service.load_optional(None);
let merged_in = merge(&local.databases, &self.dashboard_cache);
let ping_result = self.status_service.ping(&merged_in.databases).await?;
self.dashboard_cache = collect_configs(&ping_result);
if let Err(e) = persist_cache(&self.cache_path, &self.dashboard_cache) {
error!("Failed to persist dashboard cache: {e}");
}
let merged = merge(&local.databases, &self.dashboard_cache);
for db in ping_result.databases.iter() {
let database = config
let Some(database) = merged
.databases
.iter()
.find(|cfg_db| cfg_db.generated_id == db.generated_id)
.unwrap();
else {
warn!("No config for returned database {}; skipping", db.generated_id);
continue;
};
info!(
"Generated Id: {} | backup action: {} | restore action: {} | Database Name: {}",
db.generated_id, db.data.backup.action, db.data.restore.action, database.name,
@@ -59,14 +81,14 @@ impl Agent {
.backup_service
.dispatch(
&db.generated_id,
&config,
&merged,
method.clone(),
&db.storages,
db.encrypt,
)
.await;
} else if db.data.restore.action {
let _ = self.restore_service.dispatch(db, &config).await;
let _ = self.restore_service.dispatch(db, &merged).await;
}
}
+1 -1
View File
@@ -15,7 +15,7 @@ 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")
Docker::connect_with_defaults().context("Failed to connect to Docker daemon socket")
}
pub fn parse_container_id(mountinfo: &str, cgroup: &str) -> Option<String> {
+1
View File
@@ -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";
+2 -1
View File
@@ -12,7 +12,8 @@ pub async fn run(cfg: DatabaseConfig, env: HashMap<String, String>) -> anyhow::R
.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 -1
View File
@@ -1,6 +1,6 @@
pub mod docker_volume;
pub mod factory;
mod mongodb;
pub mod mongodb;
pub mod mysql;
pub mod postgres;
mod redis;
+68 -11
View File
@@ -1,6 +1,13 @@
use crate::services::config::DatabaseConfig;
use anyhow::Result;
use mongodb::Client;
use percent_encoding::{AsciiSet, NON_ALPHANUMERIC, utf8_percent_encode};
const USERINFO_ENCODE: &AsciiSet = &NON_ALPHANUMERIC
.remove(b'-')
.remove(b'_')
.remove(b'.')
.remove(b'~');
pub async fn connect(cfg: DatabaseConfig) -> Result<Client> {
let uri = get_mongo_uri(cfg)?;
@@ -16,19 +23,69 @@ pub fn select_mongo_path() -> std::path::PathBuf {
}
pub fn get_mongo_uri(cfg: DatabaseConfig) -> Result<String> {
if cfg.username.is_empty() || cfg.password.is_empty() {
Ok(format!(
"mongodb://{}:{}/{}",
cfg.host, cfg.port, cfg.database
))
} else {
Ok(format!(
"mongodb://{}:{}@{}:{}/{}?authSource=admin",
cfg.username, cfg.password, cfg.host, cfg.port, cfg.database
))
}
Ok(build_mongo_uri(&cfg, true))
}
pub fn build_mongo_uri(cfg: &DatabaseConfig, include_db: bool) -> String {
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();
let credentials = if has_auth {
format!(
"{}:{}@",
utf8_percent_encode(&cfg.username, USERINFO_ENCODE),
utf8_percent_encode(&cfg.password, USERINFO_ENCODE)
)
} else {
String::new()
};
let authority = if is_srv || is_multi_host {
cfg.host.clone()
} else {
format!("{}:{}", cfg.host, cfg.port)
};
let path = if include_db {
format!("/{}", cfg.database)
} else {
"/".to_string()
};
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)
}
pub fn extract_db_name(dry_output: &str) -> Option<String> {
let mut dbs = std::collections::HashSet::new();
+1 -1
View File
@@ -1,5 +1,5 @@
mod backup;
mod connection;
pub mod connection;
pub mod database;
mod ping;
mod restore;
+5 -1
View File
@@ -19,7 +19,11 @@ pub async fn run(cfg: DatabaseConfig) -> Result<bool> {
Ok(_) => Ok(true),
Err(e) => {
error!("--- MongoDB Connection Error Details ---");
error!("Target Host: {}:{}", cfg.host, cfg.port);
if cfg.port == 0 {
error!("Target Host: {} (srv)", cfg.host);
} else {
error!("Target Host: {}:{}", cfg.host, cfg.port);
}
error!("Error Kind: {:?}", e.kind);
error!("Full Error: {}", e);
error!("Check you database network connectivity");
+4 -8
View File
@@ -1,4 +1,6 @@
use crate::domain::mongodb::connection::{extract_db_name, get_mongo_uri, select_mongo_path};
use crate::domain::mongodb::connection::{
build_mongo_uri, extract_db_name, get_mongo_uri, select_mongo_path,
};
use crate::services::backup::logger::JobLogger;
use crate::services::config::DatabaseConfig;
use anyhow::{Context, Result};
@@ -16,13 +18,7 @@ pub async fn run(cfg: DatabaseConfig, restore_file: PathBuf, logger: Arc<JobLogg
let dry_start = Instant::now();
let dry_run = Command::new(&mongorestore)
.arg(format!(
"--uri={}",
format!(
"mongodb://{}:{}@{}:{}/?authSource=admin",
cfg.username, cfg.password, cfg.host, cfg.port
)
))
.arg(format!("--uri={}", build_mongo_uri(&cfg, false)))
.arg(format!("--archive={}", restore_file.display()))
.arg("--gzip")
.arg("--dryRun")
+2 -1
View File
@@ -12,7 +12,8 @@ pub async fn run(cfg: DatabaseConfig, env: HashMap<String, String>) -> anyhow::R
.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;
+2
View File
@@ -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;
+1
View File
@@ -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
View File
@@ -22,7 +22,6 @@ 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)"),
+8
View File
@@ -1,5 +1,6 @@
#![allow(dead_code)]
use crate::services::config::DatabaseConfig;
use crate::utils::deserializer::{deserialize_snake_case, string_or_number_to_string};
use serde::{Deserialize, Serialize};
use toml::Value;
@@ -39,6 +40,13 @@ pub struct DatabaseStatus {
pub storages_encrypted: Option<bool>,
#[serde(default)]
pub storages_ciphertext: Option<String>,
#[serde(default)]
pub config_encrypted: Option<bool>,
#[serde(default)]
pub config_ciphertext: Option<String>,
/// Filled in memory after decrypting `config_ciphertext`; never on the wire.
#[serde(skip)]
pub resolved_config: Option<DatabaseConfig>,
pub encrypt: bool,
pub data: DatabaseData,
}
+7
View File
@@ -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"))
}
}
+26 -1
View File
@@ -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,
+38 -11
View File
@@ -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(
+136 -128
View File
@@ -1,7 +1,7 @@
#![allow(dead_code)]
use crate::core::context::Context;
use serde::Deserialize;
use serde::{Deserialize, Serialize};
use serde_json;
use std::collections::HashMap;
use std::fs::File;
@@ -12,7 +12,7 @@ use toml;
use tracing::info;
use uuid::Uuid;
#[derive(Debug, Deserialize, Clone)]
#[derive(Debug, Serialize, Deserialize, Clone)]
#[serde(rename_all = "lowercase")]
pub enum DbType {
Mysql,
@@ -49,7 +49,7 @@ impl DbType {
}
#[allow(dead_code)]
#[derive(Debug, Deserialize, Clone)]
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct DatabaseConfig {
pub name: String,
pub database: String,
@@ -68,7 +68,7 @@ pub struct DatabaseConfig {
}
#[allow(dead_code)]
#[derive(Debug, Deserialize, Clone)]
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct DatabasesConfig {
pub databases: Vec<DatabaseConfig>,
}
@@ -98,6 +98,106 @@ pub struct InputDatabasesConfig {
pub databases: Vec<InputDatabaseConfig>,
}
fn required<T: Clone>(opt: &Option<T>, db_name: &str, field_name: &str) -> Result<T, String> {
match opt {
Some(v) => Ok(v.clone()),
None => Err(format!(
"Missing required field '{}' for database '{}'",
field_name, db_name
)),
}
}
fn optional<T: Clone + Default>(opt: &Option<T>) -> T {
opt.clone().unwrap_or_default()
}
pub fn build_config(db: InputDatabaseConfig) -> Result<DatabaseConfig, String> {
if Uuid::parse_str(&db.generated_id).is_err() {
return Err(format!("Invalid UUID for database '{}'", db.name));
}
let username = match db.db_type {
DbType::Postgresql
| DbType::PostgresqlCluster
| DbType::Mysql
| DbType::Mariadb
| DbType::Mssql => required(&db.username, &db.name, "username")?,
_ => optional(&db.username),
};
let password = match db.db_type {
DbType::Postgresql
| DbType::PostgresqlCluster
| DbType::Mysql
| DbType::Mariadb
| DbType::Mssql => required(&db.password, &db.name, "password")?,
_ => optional(&db.password),
};
let host = match db.db_type {
DbType::Postgresql
| DbType::PostgresqlCluster
| DbType::Mysql
| DbType::Mariadb
| DbType::MongoDB
| DbType::Redis
| DbType::Firebird
| DbType::Valkey
| DbType::Mssql => required(&db.host, &db.name, "host")?,
DbType::Sqlite | DbType::DockerVolume => optional(&db.host),
};
let port = match db.db_type {
DbType::Postgresql
| DbType::PostgresqlCluster
| DbType::Mysql
| DbType::Mariadb
| DbType::Redis
| DbType::Firebird
| DbType::Valkey
| DbType::Mssql => required(&db.port, &db.name, "port")?,
DbType::MongoDB | DbType::Sqlite | DbType::DockerVolume => db.port.unwrap_or(0),
};
let database_name = match db.db_type {
DbType::Sqlite | DbType::Redis | DbType::Valkey | DbType::DockerVolume => {
optional(&db.database)
}
DbType::PostgresqlCluster => db
.database
.clone()
.unwrap_or_else(|| "postgres".to_string()),
_ => required(&db.database, &db.name, "database")?,
};
let path_val = match db.db_type {
DbType::Sqlite => required(&db.path, &db.name, "path")?,
_ => optional(&db.path),
};
let max_packet_size = match db.db_type {
DbType::Mysql | DbType::Mariadb => db.max_packet_size.unwrap_or_else(|| "512M".to_string()),
_ => String::new(),
};
let volume_name = match db.db_type {
DbType::DockerVolume => required(&db.volume_name, &db.name, "volume_name")?,
_ => optional(&db.volume_name),
};
Ok(DatabaseConfig {
name: db.name,
database: database_name,
db_type: db.db_type,
username,
password,
host,
port,
generated_id: db.generated_id,
path: path_val,
max_packet_size,
volume_name,
container_name: db.container_name.clone(),
options: db.options.unwrap_or_default(),
})
}
pub struct ConfigService {
ctx: Arc<Context>,
}
@@ -107,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);
@@ -150,128 +253,33 @@ impl ConfigService {
_ => return Err("Unsupported config file format. Use .json or .toml".to_string()),
};
fn required<T: Clone>(
opt: &Option<T>,
db_name: &str,
field_name: &str,
) -> Result<T, String> {
match opt {
Some(v) => Ok(v.clone()),
None => {
let msg = format!(
"Missing required field '{}' for database '{}'",
field_name, db_name
);
Err(msg)
}
}
}
fn optional<T: Clone>(opt: &Option<T>) -> T
where
T: Default,
{
opt.clone().unwrap_or_default()
}
let mut databases = Vec::with_capacity(input_config.databases.len());
for db in input_config.databases {
if Uuid::parse_str(&db.generated_id).is_err() {
return Err(format!("Invalid UUID for database '{}'", db.name));
}
let username = match db.db_type {
DbType::Postgresql
| DbType::PostgresqlCluster
| DbType::Mysql
| DbType::Mariadb
| DbType::Mssql => required(&db.username, &db.name, "username")?,
_ => optional(&db.username),
};
let password = match db.db_type {
DbType::Postgresql
| DbType::PostgresqlCluster
| DbType::Mysql
| DbType::Mariadb
| DbType::Mssql => required(&db.password, &db.name, "password")?,
_ => optional(&db.password),
};
let host = match db.db_type {
DbType::Postgresql
| DbType::PostgresqlCluster
| DbType::Mysql
| DbType::Mariadb
| DbType::MongoDB
| DbType::Redis
| DbType::Firebird
| DbType::Valkey
| DbType::Mssql => required(&db.host, &db.name, "host")?,
DbType::Sqlite | DbType::DockerVolume => optional(&db.host),
};
let port = match db.db_type {
DbType::Postgresql
| DbType::PostgresqlCluster
| DbType::Mysql
| DbType::Mariadb
| DbType::MongoDB
| DbType::Redis
| DbType::Firebird
| DbType::Valkey
| DbType::Mssql => required(&db.port, &db.name, "port")?,
DbType::Sqlite | DbType::DockerVolume => db.port.unwrap_or(0),
};
let database_name = match db.db_type {
DbType::Sqlite | DbType::Redis | DbType::Valkey | DbType::DockerVolume => {
optional(&db.database)
}
DbType::PostgresqlCluster => db
.database
.clone()
.unwrap_or_else(|| "postgres".to_string()),
_ => required(&db.database, &db.name, "database")?,
};
let path_val = match db.db_type {
DbType::Sqlite => required(&db.path, &db.name, "path")?,
_ => optional(&db.path),
};
let max_packet_size = match db.db_type {
DbType::Mysql | DbType::Mariadb => {
db.max_packet_size.unwrap_or_else(|| "512M".to_string())
}
_ => 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,
db_type: db.db_type,
username,
password,
host,
port,
generated_id: db.generated_id,
path: path_val,
max_packet_size,
volume_name,
container_name,
options: db.options.unwrap_or_default(),
});
databases.push(build_config(db)?);
}
info!("Databases: {} instances loaded", databases.len());
Ok(DatabasesConfig { databases })
}
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"
);
DatabasesConfig { databases: Vec::new() }
})
}
}
+56
View File
@@ -0,0 +1,56 @@
#![allow(dead_code)]
use crate::services::api::models::agent::status::PingResult;
use crate::services::config::{DatabaseConfig, DatabasesConfig};
use std::path::Path;
pub fn merge(local: &[DatabaseConfig], dashboard: &[DatabaseConfig]) -> DatabasesConfig {
let mut databases: Vec<DatabaseConfig> = local.to_vec();
for d in dashboard {
if let Some(slot) = databases
.iter_mut()
.find(|c| c.generated_id == d.generated_id)
{
*slot = d.clone();
} else {
databases.push(d.clone());
}
}
DatabasesConfig { databases }
}
pub fn collect_configs(ping: &PingResult) -> Vec<DatabaseConfig> {
ping.databases
.iter()
.filter_map(|db| db.resolved_config.clone())
.collect()
}
pub fn load_cache(path: &Path) -> Vec<DatabaseConfig> {
let contents = match std::fs::read_to_string(path) {
Ok(c) => c,
Err(_) => return Vec::new(),
};
match serde_json::from_str::<DatabasesConfig>(&contents) {
Ok(cfg) => cfg.databases,
Err(e) => {
tracing::warn!("Dashboard cache at {:?} is corrupt ({e}); ignoring", path);
Vec::new()
}
}
}
pub fn persist_cache(path: &Path, databases: &[DatabaseConfig]) -> std::io::Result<()> {
let wrapper = DatabasesConfig {
databases: databases.to_vec(),
};
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)?;
Ok(())
}
+1
View File
@@ -2,6 +2,7 @@ pub mod api;
pub mod backup;
pub mod config;
pub mod cron;
pub mod dashboard_config;
pub mod restore;
pub mod status;
pub mod storage;
+43 -2
View File
@@ -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!(
+26 -1
View File
@@ -3,9 +3,10 @@
use crate::core::context::Context;
use crate::domain::factory::DatabaseFactory;
use crate::services::api::endpoints::status::DatabasePayload;
use crate::services::api::models::agent::status::DatabaseStatus;
use crate::services::api::models::agent::status::DatabaseStorage;
use crate::services::api::models::agent::status::PingResult;
use crate::services::config::DatabaseConfig;
use crate::services::config::{build_config, DatabaseConfig, InputDatabaseConfig};
use crate::settings::CONFIG;
use crate::utils::file::decrypt_json_gcm;
use futures_util::future::try_join_all;
@@ -14,6 +15,26 @@ use std::error::Error;
use std::sync::Arc;
use tracing::info;
pub fn resolve_dashboard_config(
status: &mut DatabaseStatus,
master_key_b64: &str,
) -> Result<(), String> {
if status.config_encrypted != Some(true) {
return Ok(());
}
let ciphertext = status
.config_ciphertext
.as_deref()
.ok_or("config_encrypted set but config_ciphertext missing")?;
let plaintext = decrypt_json_gcm(ciphertext, master_key_b64)
.map_err(|e| format!("Failed to decrypt config: {e}"))?;
let input: InputDatabaseConfig = serde_json::from_slice(&plaintext)
.map_err(|e| format!("Failed to parse decrypted config: {e}"))?;
status.resolved_config = Some(build_config(input)?);
Ok(())
}
pub struct StatusService {
ctx: Arc<Context>,
client: Client,
@@ -67,6 +88,10 @@ impl StatusService {
db.storages = serde_json::from_slice::<Vec<DatabaseStorage>>(&plaintext)
.map_err(|e| format!("Failed to parse decrypted storages: {e}"))?;
}
if let Err(e) = resolve_dashboard_config(db, &edge_key.master_key_b64) {
tracing::warn!("Skipping dashboard config for {}: {e}", db.generated_id);
}
}
Ok(result)
}
+4 -1
View File
@@ -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
View File
@@ -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 &sections {
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,
}
+12 -6
View File
@@ -13,7 +13,9 @@ use aws_config::retry::RetryConfig;
use aws_sdk_s3 as s3;
use aws_sdk_s3::config::BehaviorVersion;
use aws_sdk_s3::config::Region;
use aws_sdk_s3::config::RequestChecksumCalculation;
use aws_sdk_s3::config::retry::ReconnectMode;
use aws_sdk_s3::error::DisplayErrorContext;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart};
use futures::StreamExt;
@@ -135,6 +137,7 @@ impl StorageProvider for S3Provider {
.credentials_provider(credentials)
.region(region)
.force_path_style(true)
.request_checksum_calculation(RequestChecksumCalculation::WhenRequired)
.endpoint_url(endpoint)
.behavior_version(BehaviorVersion::latest())
.build();
@@ -164,11 +167,12 @@ impl StorageProvider for S3Provider {
{
Ok(r) => r,
Err(e) => {
error!("Failed to create multipart upload: {}", e);
let detail = DisplayErrorContext(&e).to_string();
error!("Failed to create multipart upload: {}", detail);
return UploadResult {
storage_id: storage.id.clone(),
success: false,
error: Some(e.to_string()),
error: Some(detail),
remote_file_path: None,
total_size: None,
};
@@ -251,7 +255,8 @@ impl StorageProvider for S3Provider {
}
}
Err(e) => {
error!("Failed to upload part {}: {}", part_number, e);
let detail = DisplayErrorContext(&e).to_string();
error!("Failed to upload part {}: {}", part_number, detail);
let _ = client
.abort_multipart_upload()
.bucket(bucket)
@@ -262,7 +267,7 @@ impl StorageProvider for S3Provider {
return UploadResult {
storage_id: storage.id.clone(),
success: false,
error: Some(e.to_string()),
error: Some(detail),
remote_file_path: None,
total_size: None,
};
@@ -317,7 +322,8 @@ impl StorageProvider for S3Provider {
}
}
Err(e) => {
error!("Failed to complete multipart upload: {}", e);
let detail = DisplayErrorContext(&e).to_string();
error!("Failed to complete multipart upload: {}", detail);
let _ = client
.abort_multipart_upload()
.bucket(bucket)
@@ -328,7 +334,7 @@ impl StorageProvider for S3Provider {
UploadResult {
storage_id: storage.id.clone(),
success: false,
error: Some(e.to_string()),
error: Some(detail),
remote_file_path: None,
total_size: None,
}
@@ -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))
}
+114
View 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
View File
@@ -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,
}
}
}
+99
View File
@@ -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"
);
}
+95
View File
@@ -148,3 +148,98 @@ fn database_status_encrypted_envelope() {
assert_eq!(status.storages_encrypted, Some(true));
assert_eq!(status.storages_ciphertext.as_deref(), Some("AQIDBA=="));
}
#[test]
fn database_status_defaults_config_fields_absent() {
let json = r#"{
"dbms": "postgresql",
"generatedId": "16678159-ff7e-4c97-8c83-0adeff214681",
"encrypt": false,
"data": { "backup": { "action": false, "cron": null },
"restore": { "action": false, "file": null, "metaFile": null, "size": null } }
}"#;
let status: crate::services::api::models::agent::status::DatabaseStatus =
serde_json::from_str(json).unwrap();
assert_eq!(status.config_encrypted, None);
assert!(status.config_ciphertext.is_none());
assert!(status.resolved_config.is_none());
}
#[test]
fn resolve_dashboard_config_decrypts_full_entry() {
use crate::services::status::resolve_dashboard_config;
use base64::{engine::general_purpose, Engine};
// 32-byte master key, base64 STANDARD (matches decrypt_json_gcm).
let master_key_b64 = general_purpose::STANDARD.encode([7u8; 32]);
// Full agent-entry shape the dashboard encrypts.
let entry = r#"{
"name": "Dashboard PG",
"type": "postgresql",
"database": "app",
"username": "postgres",
"password": "s3cret",
"port": 5432,
"host": "10.0.0.10",
"generated_id": "16678159-ff7e-4c97-8c83-0adeff214681"
}"#;
let ciphertext = encrypt_json_gcm(entry.as_bytes(), &master_key_b64);
let mut status: crate::services::api::models::agent::status::DatabaseStatus =
serde_json::from_str(
r#"{
"dbms": "postgresql",
"generatedId": "16678159-ff7e-4c97-8c83-0adeff214681",
"encrypt": false,
"config_encrypted": true,
"config_ciphertext": "PLACEHOLDER",
"data": { "backup": { "action": false, "cron": null },
"restore": { "action": false, "file": null, "metaFile": null, "size": null } }
}"#,
)
.unwrap();
status.config_ciphertext = Some(ciphertext);
resolve_dashboard_config(&mut status, &master_key_b64).unwrap();
let cfg = status.resolved_config.expect("resolved");
assert_eq!(cfg.name, "Dashboard PG");
assert_eq!(cfg.password, "s3cret");
assert_eq!(cfg.host, "10.0.0.10");
assert_eq!(cfg.db_type.as_str(), "postgresql");
}
#[test]
fn resolve_dashboard_config_noop_when_not_encrypted() {
use crate::services::status::resolve_dashboard_config;
let mut status: crate::services::api::models::agent::status::DatabaseStatus =
serde_json::from_str(
r#"{
"dbms": "postgresql",
"generatedId": "16678159-ff7e-4c97-8c83-0adeff214681",
"encrypt": false,
"data": { "backup": { "action": false, "cron": null },
"restore": { "action": false, "file": null, "metaFile": null, "size": null } }
}"#,
)
.unwrap();
resolve_dashboard_config(&mut status, "unused").unwrap();
assert!(status.resolved_config.is_none());
}
fn encrypt_json_gcm(plaintext: &[u8], master_key_b64: &str) -> String {
use aes_gcm::aead::{Aead, KeyInit};
use aes_gcm::{Aes256Gcm, Key, Nonce};
use base64::{engine::general_purpose, Engine};
let key_bytes = general_purpose::STANDARD.decode(master_key_b64).unwrap();
let key = Key::<Aes256Gcm>::try_from(key_bytes.as_slice()).unwrap();
let cipher = Aes256Gcm::new(&key);
let nonce_bytes = [0u8; 12];
let nonce = Nonce::try_from(&nonce_bytes[..]).unwrap();
let ct = cipher.encrypt(&nonce, plaintext).unwrap();
let mut data = nonce_bytes.to_vec();
data.extend_from_slice(&ct);
general_purpose::STANDARD.encode(data)
}
+68
View File
@@ -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<_>>()
);
}
+110
View File
@@ -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")
);
}
+71 -4
View File
@@ -1,15 +1,12 @@
use crate::core::context::Context;
use crate::services::api::ApiClient;
use crate::services::config::ConfigService;
use crate::services::config::{build_config, DatabasesConfig, InputDatabaseConfig};
use crate::utils::edge_key::EdgeKey;
use std::io::Write;
use std::sync::Arc;
use tempfile::NamedTempFile;
// `ConfigService::load` never touches `self.ctx` on the `Some(file_path)` path,
// so the values here don't matter — but `Context::new()` panics without an
// `EDGE_KEY` env var, so build the struct directly (mirrors
// backup_uploader_tests.rs's `ctx_pointing_at`).
fn test_context() -> Arc<Context> {
Arc::new(Context {
edge_key: EdgeKey {
@@ -264,3 +261,73 @@ fn docker_volume_requires_volume_name() {
let err = service.load(Some(file.path().to_str().unwrap())).unwrap_err();
assert!(err.contains("volume_name"), "error was: {err}");
}
#[test]
fn build_config_applies_type_defaults() {
let input: InputDatabaseConfig = serde_json::from_str(
r#"{
"name": "cluster1",
"type": "postgresql-cluster",
"username": "postgres",
"password": "p",
"port": 5432,
"host": "localhost",
"generated_id": "16678159-ff7e-4c97-8c83-0adeff214681"
}"#,
)
.unwrap();
let cfg = build_config(input).unwrap();
assert_eq!(cfg.db_type.as_str(), "postgresql-cluster");
assert_eq!(cfg.database, "postgres"); // cluster default
}
#[test]
fn build_config_rejects_missing_required_field() {
let input: InputDatabaseConfig = serde_json::from_str(
r#"{
"name": "pg",
"type": "postgresql",
"username": "postgres",
"port": 5432,
"host": "localhost",
"generated_id": "16678159-ff7e-4c97-8c83-0adeff214681"
}"#,
)
.unwrap();
let err = build_config(input).unwrap_err();
assert!(err.contains("password"), "unexpected error: {err}");
}
#[test]
fn load_optional_returns_empty_when_file_missing() {
let service = ConfigService::new(test_context());
let cfg = service.load_optional(Some("/nonexistent/path/does-not-exist.json"));
assert!(cfg.databases.is_empty());
}
#[test]
fn databases_config_roundtrips_through_serde() {
let input: InputDatabaseConfig = serde_json::from_str(
r#"{
"name": "pg",
"type": "postgresql",
"database": "app",
"username": "postgres",
"password": "secret",
"port": 5432,
"host": "localhost",
"generated_id": "16678159-ff7e-4c97-8c83-0adeff214681"
}"#,
)
.unwrap();
let cfg = build_config(input).unwrap();
let wrapped = DatabasesConfig { databases: vec![cfg] };
let json = serde_json::to_string(&wrapped).unwrap();
let back: DatabasesConfig = serde_json::from_str(&json).unwrap();
assert_eq!(back.databases[0].name, "pg");
assert_eq!(back.databases[0].db_type.as_str(), "postgresql");
assert_eq!(back.databases[0].password, "secret");
}
@@ -0,0 +1,84 @@
use crate::services::config::{build_config, DatabaseConfig, InputDatabaseConfig};
use crate::services::dashboard_config::merge;
use crate::services::dashboard_config::{load_cache, persist_cache};
fn cfg(name: &str, gen_id: &str, host: &str) -> DatabaseConfig {
let json = format!(
r#"{{ "name": "{name}", "type": "postgresql", "database": "app",
"username": "u", "password": "p", "port": 5432,
"host": "{host}", "generated_id": "{gen_id}" }}"#
);
let input: InputDatabaseConfig = serde_json::from_str(&json).unwrap();
build_config(input).unwrap()
}
const ID_A: &str = "16678159-ff7e-4c97-8c83-0adeff214681";
const ID_B: &str = "16678124-ff7e-4c97-8c83-0adeff214681";
#[test]
fn merge_keeps_local_only_databases() {
let local = vec![cfg("local-a", ID_A, "local-host")];
let merged = merge(&local, &[]);
assert_eq!(merged.databases.len(), 1);
assert_eq!(merged.databases[0].host, "local-host");
}
#[test]
fn merge_appends_dashboard_only_databases() {
let local = vec![cfg("local-a", ID_A, "local-host")];
let dashboard = vec![cfg("dash-b", ID_B, "dash-host")];
let merged = merge(&local, &dashboard);
assert_eq!(merged.databases.len(), 2);
assert!(merged.databases.iter().any(|d| d.generated_id == ID_B));
}
#[test]
fn merge_dashboard_wins_on_id_collision() {
let local = vec![cfg("local-a", ID_A, "local-host")];
let dashboard = vec![cfg("dash-a", ID_A, "dash-host")];
let merged = merge(&local, &dashboard);
assert_eq!(merged.databases.len(), 1);
assert_eq!(merged.databases[0].host, "dash-host"); // dashboard wins
assert_eq!(merged.databases[0].name, "dash-a");
}
#[test]
fn cache_roundtrips() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("dashboard_databases.json");
let dbs = vec![cfg("dash-a", ID_A, "dash-host")];
persist_cache(&path, &dbs).unwrap();
let loaded = load_cache(&path);
assert_eq!(loaded.len(), 1);
assert_eq!(loaded[0].generated_id, ID_A);
assert_eq!(loaded[0].host, "dash-host");
}
#[test]
fn load_cache_missing_file_is_empty() {
let loaded = load_cache(std::path::Path::new("/nonexistent/dashboard_databases.json"));
assert!(loaded.is_empty());
}
#[test]
fn load_cache_corrupt_file_is_empty() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("dashboard_databases.json");
std::fs::write(&path, b"{ this is not valid json").unwrap();
let loaded = load_cache(&path);
assert!(loaded.is_empty());
}
#[test]
fn persist_cache_leaves_no_tmp_file() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("dashboard_databases.json");
persist_cache(&path, &[cfg("dash-a", ID_A, "h")]).unwrap();
let tmp = path.with_extension("json.tmp");
assert!(!tmp.exists(), "temp file should have been renamed away");
assert!(path.exists());
}
+3
View File
@@ -1,3 +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")
);
}
+2
View File
@@ -1,2 +1,4 @@
mod azure_blob;
mod google_cloud_storage;
mod rclone;
mod sftp;
+455
View File
@@ -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"
);
}
+75
View File
@@ -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());
}
+1
View File
@@ -4,4 +4,5 @@ mod deserializer;
mod edge_key_tests;
mod file_tests;
mod normalize_cron_tests;
mod retry_tests;
mod stream_tests;
+135
View File
@@ -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));
}
+1
View File
@@ -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;
+75
View File
@@ -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);
}
}
}
}
+6 -1
View File
@@ -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()