mirror of
https://github.com/Portabase/agent.git
synced 2026-09-11 02:27:10 +00:00
Compare commits
28 Commits
1.1.6-rc.1
..
1.2.5
| Author | SHA1 | Date | |
|---|---|---|---|
| 189ad866de | |||
| 55bbd4e724 | |||
| 4f67a57681 | |||
| 45a1118f64 | |||
| f2275d5aca | |||
| b136140b55 | |||
| 773077f66b | |||
| 46ff086cf9 | |||
| f563c187db | |||
| df0dffa47b | |||
| 864396422b | |||
| e280c1d2f8 | |||
| 91a9d0654b | |||
| bc857a4c18 | |||
| 2454915866 | |||
| a4a7a9a2c4 | |||
| 997ba2afa2 | |||
| f973bc1f52 | |||
| 6ea9fdf8bc | |||
| ad68be8a5a | |||
| 434c131afb | |||
| 675ac5976d | |||
| 70543fe032 | |||
| 77f90308b4 | |||
| 63cb254abd | |||
| 7738745def | |||
| 82659f681c | |||
| eee3bebd50 |
@@ -0,0 +1,23 @@
|
||||
reviews:
|
||||
auto_review:
|
||||
enabled: true
|
||||
ignore_usernames: ["dependabot[bot]"]
|
||||
labels: ["!wip", "!draft"]
|
||||
auto_apply_labels: true
|
||||
suggested_reviewers: true
|
||||
auto_assign_reviewers: true
|
||||
|
||||
issue_enrichment:
|
||||
auto_enrich:
|
||||
enabled: true
|
||||
labeling:
|
||||
auto_apply_labels: true
|
||||
labeling_instructions:
|
||||
- label: bug
|
||||
instructions: "Error, crash, or incorrect behavior reports"
|
||||
- label: enhancement
|
||||
instructions: "Feature requests or improvements"
|
||||
- label: documentation
|
||||
instructions: "Docs updates or missing documentation"
|
||||
- label: good first issue
|
||||
instructions: "Beginner-friendly tasks suitable for new contributors"
|
||||
@@ -0,0 +1,73 @@
|
||||
name: Discord Notification
|
||||
|
||||
on:
|
||||
workflow_call:
|
||||
inputs:
|
||||
release_tag:
|
||||
required: true
|
||||
type: string
|
||||
discord_title:
|
||||
required: true
|
||||
type: string
|
||||
discord_color:
|
||||
required: true
|
||||
type: number
|
||||
discord_footer:
|
||||
required: true
|
||||
type: string
|
||||
secrets:
|
||||
DISCORD_WEBHOOK:
|
||||
required: true
|
||||
GH_TOKEN:
|
||||
required: true
|
||||
|
||||
jobs:
|
||||
notify-discord:
|
||||
runs-on: ubuntu-latest
|
||||
steps:
|
||||
- name: Send Discord Notification
|
||||
env:
|
||||
DISCORD_WEBHOOK: ${{ secrets.DISCORD_WEBHOOK }}
|
||||
GH_TOKEN: ${{ secrets.GH_TOKEN }}
|
||||
run: |
|
||||
RELEASE_INFO=$(gh release view "${{ inputs.release_tag }}" -R ${{ github.repository }} --json name,url,body,author)
|
||||
|
||||
RELEASE_TITLE=$(echo "$RELEASE_INFO" | jq -r .name)
|
||||
if [ -z "$RELEASE_TITLE" ] || [ "$RELEASE_TITLE" = "null" ]; then RELEASE_TITLE="${{ inputs.release_tag }}"; fi
|
||||
|
||||
RELEASE_URL=$(echo "$RELEASE_INFO" | jq -r .url)
|
||||
RELEASE_BODY=$(echo "$RELEASE_INFO" | jq -r .body)
|
||||
|
||||
AUTHOR_NAME="Portabase"
|
||||
AUTHOR_ICON="https://github.com/Portabase.png"
|
||||
|
||||
PAYLOAD=$(jq -n \
|
||||
--arg title "$RELEASE_TITLE" \
|
||||
--arg description "$RELEASE_BODY" \
|
||||
--arg url "$RELEASE_URL" \
|
||||
--arg author "$AUTHOR_NAME" \
|
||||
--arg icon "$AUTHOR_ICON" \
|
||||
--arg discord_title "${{ inputs.discord_title }}" \
|
||||
--arg discord_footer "${{ inputs.discord_footer }}" \
|
||||
--argjson discord_color ${{ inputs.discord_color }} \
|
||||
'{
|
||||
content: $discord_title,
|
||||
embeds: [{
|
||||
title: $title,
|
||||
url: $url,
|
||||
description: $description,
|
||||
color: $discord_color,
|
||||
author: {
|
||||
name: $author,
|
||||
icon_url: $icon
|
||||
},
|
||||
footer: {
|
||||
text: $discord_footer
|
||||
}
|
||||
}]
|
||||
}'
|
||||
)
|
||||
|
||||
curl -H "Content-Type: application/json" \
|
||||
-d "$PAYLOAD" \
|
||||
"$DISCORD_WEBHOOK"
|
||||
@@ -3,10 +3,16 @@ name: Docker Publish
|
||||
on:
|
||||
workflow_call:
|
||||
inputs:
|
||||
version:
|
||||
required: true
|
||||
type: string
|
||||
ref:
|
||||
required: true
|
||||
type: string
|
||||
image_name:
|
||||
required: false
|
||||
type: string
|
||||
default: 'portabase/agent'
|
||||
default: "portabase/agent"
|
||||
add_latest:
|
||||
required: false
|
||||
type: boolean
|
||||
@@ -14,11 +20,11 @@ on:
|
||||
target:
|
||||
required: false
|
||||
type: string
|
||||
default: 'prod'
|
||||
default: "prod"
|
||||
dockerfile:
|
||||
required: false
|
||||
type: string
|
||||
default: './docker/Dockerfile'
|
||||
default: "./docker/Dockerfile"
|
||||
secrets:
|
||||
DOCKER_USERNAME:
|
||||
required: true
|
||||
@@ -27,32 +33,34 @@ on:
|
||||
|
||||
jobs:
|
||||
build:
|
||||
name: Build and push Docker images
|
||||
runs-on: ${{ matrix.platform == 'linux/amd64' && 'ubuntu-latest' || matrix.platform == 'linux/arm64' && 'ubuntu-24.04-arm' }}
|
||||
name: Build architectures
|
||||
runs-on: ${{ matrix.platform == 'linux/amd64' && 'ubuntu-latest' || 'ubuntu-24.04-arm' }}
|
||||
strategy:
|
||||
fail-fast: false
|
||||
matrix:
|
||||
platform: [linux/amd64, linux/arm64]
|
||||
platform: [ linux/amd64, linux/arm64 ]
|
||||
steps:
|
||||
- name: Checkout
|
||||
uses: actions/checkout@v4
|
||||
- uses: actions/checkout@v4
|
||||
with:
|
||||
ref: ${{ inputs.ref }}
|
||||
fetch-depth: 0
|
||||
|
||||
- name: Set up Docker Buildx
|
||||
uses: docker/setup-buildx-action@v3
|
||||
|
||||
- name: Login to Docker
|
||||
uses: docker/login-action@f4ef78c080cd8ba55a85445d5b36e214a81df20a
|
||||
uses: docker/login-action@v3
|
||||
with:
|
||||
username: ${{ secrets.DOCKER_USERNAME }}
|
||||
password: ${{ secrets.DOCKER_PASSWORD }}
|
||||
|
||||
- name: Set image tag
|
||||
id: set-tags
|
||||
- name: Prepare Image Tags
|
||||
id: prep
|
||||
run: |
|
||||
REF_NAME=${GITHUB_REF#refs/tags/}
|
||||
if [ "${{ matrix.platform }}" = "linux/amd64" ]; then
|
||||
IMAGE="${{ inputs.image_name }}:$REF_NAME-amd64"
|
||||
else
|
||||
IMAGE="${{ inputs.image_name }}:$REF_NAME-arm64"
|
||||
fi
|
||||
echo "image=$IMAGE" >> $GITHUB_OUTPUT
|
||||
ARCH=${{ matrix.platform == 'linux/amd64' && 'amd64' || 'arm64' }}
|
||||
echo "image=${{ inputs.image_name }}:${{inputs.version}}-$ARCH" >> $GITHUB_OUTPUT
|
||||
echo "safe_platform=${{ matrix.platform == 'linux/amd64' && 'linux-amd64' || 'linux-arm64' }}" >> $GITHUB_OUTPUT
|
||||
|
||||
|
||||
- name: Build and push image
|
||||
uses: docker/build-push-action@v6
|
||||
@@ -61,57 +69,71 @@ jobs:
|
||||
file: ${{ inputs.dockerfile }}
|
||||
platforms: ${{ matrix.platform }}
|
||||
push: true
|
||||
tags: ${{ steps.set-tags.outputs.image }}
|
||||
tags: ${{ steps.prep.outputs.image }}
|
||||
target: ${{ inputs.target }}
|
||||
|
||||
- name: Prepare artifact name
|
||||
id: artifact
|
||||
run: |
|
||||
platform=${{ matrix.platform }}
|
||||
echo "safe_platform=${platform//\//-}" >> $GITHUB_OUTPUT
|
||||
echo "${{ steps.set-tags.outputs.image }}" > image.txt
|
||||
- name: Save image name for manifest
|
||||
run: echo "${{ steps.prep.outputs.image }}" > image.txt
|
||||
|
||||
- uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f
|
||||
- uses: actions/upload-artifact@v4
|
||||
with:
|
||||
name: image-${{ steps.artifact.outputs.safe_platform }}
|
||||
name: image-${{ steps.prep.outputs.safe_platform }}
|
||||
path: image.txt
|
||||
if-no-files-found: warn
|
||||
compression-level: 6
|
||||
overwrite: false
|
||||
include-hidden-files: false
|
||||
retention-days: 1
|
||||
|
||||
create-manifest:
|
||||
name: Create multi-arch Docker manifest
|
||||
name: Create multi-arch manifest
|
||||
runs-on: ubuntu-latest
|
||||
needs: build
|
||||
steps:
|
||||
- uses: actions/download-artifact@37930b1c2abaa49bbe596cd826c3c89aef350131
|
||||
- uses: actions/download-artifact@v4
|
||||
with:
|
||||
name: image-linux-amd64
|
||||
path: /tmp/digests/amd64
|
||||
|
||||
- uses: actions/download-artifact@37930b1c2abaa49bbe596cd826c3c89aef350131
|
||||
- uses: actions/download-artifact@v4
|
||||
with:
|
||||
name: image-linux-arm64
|
||||
path: /tmp/digests/arm64
|
||||
|
||||
- name: Login to Docker
|
||||
uses: docker/login-action@f4ef78c080cd8ba55a85445d5b36e214a81df20a
|
||||
uses: docker/login-action@v3
|
||||
with:
|
||||
username: ${{ secrets.DOCKER_USERNAME }}
|
||||
password: ${{ secrets.DOCKER_PASSWORD }}
|
||||
|
||||
- name: Generate semantic tags
|
||||
id: tags
|
||||
run: |
|
||||
VERSION="${{ inputs.version }}"
|
||||
IFS='.' read -r MAJOR MINOR PATCH <<< "$VERSION"
|
||||
|
||||
echo "VERSION_TAG=$VERSION" >> $GITHUB_OUTPUT
|
||||
if [ -n "$MINOR" ]; then
|
||||
echo "MINOR_TAG=$MAJOR.$MINOR" >> $GITHUB_OUTPUT
|
||||
fi
|
||||
if [ -n "$MAJOR" ]; then
|
||||
echo "MAJOR_TAG=$MAJOR" >> $GITHUB_OUTPUT
|
||||
fi
|
||||
if [ "${{ inputs.add_latest }}" = "true" ]; then
|
||||
echo "LATEST_TAG=latest" >> $GITHUB_OUTPUT
|
||||
fi
|
||||
|
||||
- name: Extract Docker metadata
|
||||
id: meta
|
||||
uses: docker/metadata-action@v5
|
||||
with:
|
||||
images: ${{ inputs.image_name }}
|
||||
tags: |
|
||||
type=raw,value=${{ steps.tags.outputs.VERSION_TAG }}
|
||||
type=raw,value=${{ steps.tags.outputs.MINOR_TAG }}
|
||||
type=raw,value=${{ steps.tags.outputs.MAJOR_TAG }}
|
||||
type=raw,value=${{ steps.tags.outputs.LATEST_TAG }}
|
||||
|
||||
- name: Create and push manifest list
|
||||
working-directory: /tmp/digests
|
||||
run: |
|
||||
DOCKER_IMAGES="$(cat amd64/image.txt) $(cat arm64/image.txt)"
|
||||
REF_NAME=${GITHUB_REF#refs/tags/}
|
||||
MANIFEST_IMAGE="${{ inputs.image_name }}:$REF_NAME"
|
||||
|
||||
docker buildx imagetools create $DOCKER_IMAGES -t $MANIFEST_IMAGE
|
||||
docker buildx imagetools inspect $MANIFEST_IMAGE
|
||||
|
||||
if [ "${{ inputs.add_latest }}" = "true" ]; then
|
||||
docker buildx imagetools create $DOCKER_IMAGES -t ${{ inputs.image_name }}:latest
|
||||
docker buildx imagetools inspect ${{ inputs.image_name }}:latest
|
||||
fi
|
||||
TAG_ARGS=$(jq -cr '.tags | map("-t " + .) | join(" ")' <<< "$DOCKER_METADATA_OUTPUT_JSON")
|
||||
echo $TAG_ARGS
|
||||
docker buildx imagetools create $TAG_ARGS $DOCKER_IMAGES
|
||||
|
||||
@@ -1,128 +0,0 @@
|
||||
name: GitHub Release
|
||||
|
||||
on:
|
||||
workflow_call:
|
||||
inputs:
|
||||
prerelease:
|
||||
required: false
|
||||
type: boolean
|
||||
default: false
|
||||
make_latest:
|
||||
required: false
|
||||
type: boolean
|
||||
default: false
|
||||
discord_title:
|
||||
required: true
|
||||
type: string
|
||||
discord_color:
|
||||
required: true
|
||||
type: number
|
||||
discord_footer:
|
||||
required: true
|
||||
type: string
|
||||
secrets:
|
||||
DISCORD_WEBHOOK:
|
||||
required: true
|
||||
GH_TOKEN:
|
||||
required: true
|
||||
|
||||
jobs:
|
||||
create-release:
|
||||
runs-on: ubuntu-latest
|
||||
permissions:
|
||||
contents: write
|
||||
steps:
|
||||
- name: Check out the repo
|
||||
uses: actions/checkout@v4
|
||||
with:
|
||||
fetch-depth: 0
|
||||
|
||||
- name: Build Changelog
|
||||
id: build_changelog
|
||||
uses: mikepenz/release-changelog-builder-action@v5
|
||||
with:
|
||||
mode: "COMMIT"
|
||||
configurationJson: |
|
||||
{
|
||||
"template": "#{{CHANGELOG}}",
|
||||
"categories": [
|
||||
{
|
||||
"title": "## Feature",
|
||||
"labels": ["feat", "feature"]
|
||||
},
|
||||
{
|
||||
"title": "## Fix",
|
||||
"labels": ["fix", "bug"]
|
||||
},
|
||||
{
|
||||
"title": "## Other",
|
||||
"labels": []
|
||||
}
|
||||
],
|
||||
"label_extractor": [
|
||||
{
|
||||
"pattern": "^(build|chore|ci|docs|feat|fix|perf|refactor|revert|style|test){1}(\([\w\-\.]+\))?(!)?: ([\w ])+([\s\S]*)",
|
||||
"on_property": "title",
|
||||
"target": "$1"
|
||||
}
|
||||
]
|
||||
}
|
||||
env:
|
||||
GITHUB_TOKEN: ${{ secrets.GH_TOKEN }}
|
||||
|
||||
- name: Create GitHub Release
|
||||
uses: softprops/action-gh-release@v2
|
||||
with:
|
||||
generate_release_notes: false
|
||||
body: ${{ steps.build_changelog.outputs.changelog }}
|
||||
prerelease: ${{ inputs.prerelease }}
|
||||
make_latest: ${{ inputs.make_latest }}
|
||||
env:
|
||||
GITHUB_TOKEN: ${{ secrets.GH_TOKEN }}
|
||||
|
||||
- name: Send Discord Notification
|
||||
env:
|
||||
DISCORD_WEBHOOK: ${{ secrets.DISCORD_WEBHOOK }}
|
||||
GH_TOKEN: ${{ secrets.GH_TOKEN }}
|
||||
run: |
|
||||
RELEASE_INFO=$(gh release view "${{ github.ref_name }}" -R ${{ github.repository }} --json name,url,body,author)
|
||||
|
||||
RELEASE_TITLE=$(echo "$RELEASE_INFO" | jq -r .name)
|
||||
if [ -z "$RELEASE_TITLE" ] || [ "$RELEASE_TITLE" = "null" ]; then RELEASE_TITLE="${{ github.ref_name }}"; fi
|
||||
|
||||
RELEASE_URL=$(echo "$RELEASE_INFO" | jq -r .url)
|
||||
RELEASE_BODY=$(echo "$RELEASE_INFO" | jq -r .body)
|
||||
|
||||
AUTHOR_NAME="Portabase"
|
||||
AUTHOR_ICON="https://github.com/Portabase.png"
|
||||
|
||||
PAYLOAD=$(jq -n \
|
||||
--arg title "$RELEASE_TITLE" \
|
||||
--arg description "$RELEASE_BODY" \
|
||||
--arg url "$RELEASE_URL" \
|
||||
--arg author "$AUTHOR_NAME" \
|
||||
--arg icon "$AUTHOR_ICON" \
|
||||
--arg discord_title "${{ inputs.discord_title }}" \
|
||||
--arg discord_footer "${{ inputs.discord_footer }}" \
|
||||
--argjson discord_color ${{ inputs.discord_color }} \
|
||||
'{
|
||||
content: $discord_title,
|
||||
embeds: [{
|
||||
title: $title,
|
||||
url: $url,
|
||||
description: $description,
|
||||
color: $discord_color,
|
||||
author: {
|
||||
name: $author,
|
||||
icon_url: $icon
|
||||
},
|
||||
footer: {
|
||||
text: $discord_footer
|
||||
}
|
||||
}]
|
||||
}'
|
||||
)
|
||||
|
||||
curl -H "Content-Type: application/json" \
|
||||
-d "$PAYLOAD" \
|
||||
"$DISCORD_WEBHOOK"
|
||||
+118
-12
@@ -1,32 +1,138 @@
|
||||
name: Publish Docker image for release
|
||||
name: Auto Release & Publish
|
||||
|
||||
on:
|
||||
push:
|
||||
tags:
|
||||
- '*.*.*'
|
||||
- '!*-*'
|
||||
pull_request:
|
||||
types: [ closed ]
|
||||
branches:
|
||||
- main
|
||||
|
||||
permissions:
|
||||
contents: write
|
||||
packages: write
|
||||
|
||||
jobs:
|
||||
docker_publish:
|
||||
|
||||
check-skip:
|
||||
runs-on: ubuntu-latest
|
||||
outputs:
|
||||
skip: ${{ steps.set-skip.outputs.skip }}
|
||||
steps:
|
||||
- name: Determine if release should be skipped
|
||||
id: set-skip
|
||||
run: |
|
||||
TITLE="${{ github.event.pull_request.title }}"
|
||||
echo "PR title: $TITLE"
|
||||
|
||||
if [[ "$TITLE" == *"[skip-release]"* ]]; then
|
||||
echo "PR title contains [skip-release], skipping release jobs."
|
||||
echo "skip=true" >> $GITHUB_OUTPUT
|
||||
else
|
||||
echo "skip=false" >> $GITHUB_OUTPUT
|
||||
fi
|
||||
|
||||
create-release:
|
||||
needs: check-skip
|
||||
if: ${{ needs.check-skip.outputs.skip == 'false' }}
|
||||
runs-on: ubuntu-latest
|
||||
outputs:
|
||||
draft_tag: ${{ steps.release_step.outputs.draft_tag }}
|
||||
version: ${{ steps.release_step.outputs.version }}
|
||||
steps:
|
||||
- uses: actions/create-github-app-token@v1
|
||||
id: app-token
|
||||
with:
|
||||
app-id: ${{ vars.APP_ID }}
|
||||
private-key: ${{ secrets.APP_PRIVATE_KEY }}
|
||||
|
||||
- uses: actions/checkout@v4
|
||||
with:
|
||||
fetch-depth: 0
|
||||
token: ${{ steps.app-token.outputs.token }}
|
||||
ref: main
|
||||
|
||||
- name: Setup Rust
|
||||
uses: dtolnay/rust-toolchain@stable
|
||||
|
||||
- name: Install cargo-edit
|
||||
run: cargo install cargo-edit
|
||||
|
||||
- uses: actions/setup-node@v4
|
||||
with:
|
||||
node-version: "lts/*"
|
||||
|
||||
- name: Install release-it globally
|
||||
run: |
|
||||
npm install -g release-it
|
||||
npm install -g @release-it/conventional-changelog
|
||||
npm install -g @release-it/bumper
|
||||
|
||||
- run: |
|
||||
git config --global user.name 'github-actions[bot]'
|
||||
git config --global user.email 'github-actions[bot]@users.noreply.github.com'
|
||||
|
||||
- name: Run release-it
|
||||
id: release_step
|
||||
run: |
|
||||
git pull origin main
|
||||
|
||||
VERSION=$(release-it --ci --release-version)
|
||||
echo $VERSION
|
||||
echo "version=$VERSION" >> $GITHUB_OUTPUT
|
||||
|
||||
OUTPUT=$(release-it --ci)
|
||||
echo "$OUTPUT"
|
||||
|
||||
DRAFT_TAG=$(echo "$OUTPUT" | grep -oE 'untagged-[a-z0-9]+')
|
||||
echo $DRAFT_TAG
|
||||
echo "draft_tag=$DRAFT_TAG" >> $GITHUB_OUTPUT
|
||||
env:
|
||||
GITHUB_TOKEN: ${{ steps.app-token.outputs.token }}
|
||||
|
||||
publish-docker:
|
||||
needs: create-release
|
||||
if: ${{ needs.create-release.result == 'success' }}
|
||||
uses: ./.github/workflows/docker.yml
|
||||
with:
|
||||
version: ${{ needs.create-release.outputs.version }}
|
||||
ref: ${{ needs.create-release.outputs.version }}
|
||||
add_latest: true
|
||||
secrets:
|
||||
DOCKER_USERNAME: ${{ secrets.DOCKER_USERNAME }}
|
||||
DOCKER_PASSWORD: ${{ secrets.DOCKER_PASSWORD }}
|
||||
|
||||
github_release:
|
||||
needs: docker_publish
|
||||
uses: ./.github/workflows/github.yml
|
||||
|
||||
finalize-release:
|
||||
needs:
|
||||
- create-release
|
||||
- publish-docker
|
||||
runs-on: ubuntu-latest
|
||||
outputs:
|
||||
release_tag: ${{ steps.publish_release_step.outputs.release_tag }}
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
with:
|
||||
fetch-depth: 0
|
||||
|
||||
- name: Publish GitHub Release
|
||||
id: publish_release_step
|
||||
env:
|
||||
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
run: |
|
||||
OUTPUT=$(gh release edit ${{ needs.create-release.outputs.draft_tag }} --draft=false)
|
||||
echo "$OUTPUT"
|
||||
RELEASE_TAG=$(echo "$OUTPUT" | sed -E 's|.*/releases/tag/||')
|
||||
echo "release_tag=$RELEASE_TAG" >> $GITHUB_OUTPUT
|
||||
|
||||
notify-discord:
|
||||
needs:
|
||||
- publish-docker
|
||||
- create-release
|
||||
- finalize-release
|
||||
uses: ./.github/workflows/discord.yml
|
||||
with:
|
||||
prerelease: false
|
||||
make_latest: true
|
||||
release_tag: ${{ needs.create-release.outputs.version }}
|
||||
discord_title: "||@release-agent|| New release published"
|
||||
discord_color: 5814783
|
||||
discord_color: 3066993
|
||||
discord_footer: "Portabase"
|
||||
secrets:
|
||||
DISCORD_WEBHOOK: ${{ secrets.DISCORD_WEBHOOK }}
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
{
|
||||
"github": {
|
||||
"release": true,
|
||||
"draft": true,
|
||||
"tokenRef": "GITHUB_TOKEN"
|
||||
},
|
||||
"git": {
|
||||
"commit": true,
|
||||
"commitMessage": "chore: release ${version}",
|
||||
"requireCleanWorkingDir": true,
|
||||
"tag": true,
|
||||
"tagName": "${version}",
|
||||
"push": true
|
||||
},
|
||||
"hooks": {
|
||||
"before:bump": "cargo set-version ${version}"
|
||||
},
|
||||
"plugins": {
|
||||
"@release-it/conventional-changelog": {
|
||||
"preset": {
|
||||
"name": "conventionalcommits",
|
||||
"types": [
|
||||
{
|
||||
"type": "feat",
|
||||
"section": "✨ Features"
|
||||
},
|
||||
{
|
||||
"type": "fix",
|
||||
"section": "🐛 Bug Fixes"
|
||||
},
|
||||
{
|
||||
"type": "perf",
|
||||
"section": "⚡️ Performance Improvements"
|
||||
},
|
||||
{
|
||||
"type": "revert",
|
||||
"section": "⏪️ Reverts"
|
||||
},
|
||||
{
|
||||
"type": "docs",
|
||||
"section": "📝 Documentation",
|
||||
"hidden": true
|
||||
},
|
||||
{
|
||||
"type": "chore",
|
||||
"section": "🔧 Chores",
|
||||
"hidden": true
|
||||
}
|
||||
]
|
||||
}
|
||||
},
|
||||
"@release-it/bumper": {
|
||||
"out": [
|
||||
{
|
||||
"file": "CITATION.cff",
|
||||
"path": "version",
|
||||
"type": "text/yaml"
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
+9
-4
@@ -1,6 +1,6 @@
|
||||
cff-version: 1.2.0
|
||||
title: Portabase Agent (Rust)
|
||||
message: "If you use this software, please cite it as below."
|
||||
message: If you use this software, please cite it as below.
|
||||
type: software
|
||||
authors:
|
||||
- family-names: Gauthereau
|
||||
@@ -9,7 +9,12 @@ authors:
|
||||
given-names: Killian
|
||||
repository-code: https://github.com/Portabase/agent-rust
|
||||
url: https://portabase.io
|
||||
abstract: "Portabase is a free, open-source, self-hosted solution for database administration, providing backup and restore capabilities, scheduling, retention policies, notifications, and support for multiple storage backends. Its headless agent architecture enables connection to multiple database instances securely and efficiently."
|
||||
abstract: >-
|
||||
Portabase is a free, open-source, self-hosted solution for database
|
||||
administration, providing backup and restore capabilities, scheduling,
|
||||
retention policies, notifications, and support for multiple storage backends.
|
||||
Its headless agent architecture enables connection to multiple database
|
||||
instances securely and efficiently.
|
||||
keywords:
|
||||
- database
|
||||
- administration
|
||||
@@ -22,5 +27,5 @@ keywords:
|
||||
- self-hosted
|
||||
- portabase
|
||||
license: Apache-2.0
|
||||
version: 1.1.6-rc.1
|
||||
date-released: "2026-02-23"
|
||||
version: 1.2.5
|
||||
date-released: '2026-02-24'
|
||||
|
||||
Generated
+268
-241
File diff suppressed because it is too large
Load Diff
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "portabase-agent"
|
||||
version = "1.1.6-rc.1"
|
||||
version = "1.2.5"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
|
||||
@@ -19,6 +19,7 @@
|
||||
[](https://www.python.org/downloads/release/python-3120/)
|
||||
[](https://www.postgresql.org/)
|
||||
[](https://www.mysql.com/)
|
||||
[](https://sqlite.org/)
|
||||
[](https://mariadb.org/)
|
||||
[](https://www.mongodb.com/)
|
||||
[](https://github.com/Portabase/portabase)
|
||||
@@ -79,3 +80,5 @@ Distributed under the Apache License. See `LICENSE.txt` for more information.
|
||||
|
||||
[Docker-url]: https://www.docker.com/
|
||||
|
||||
|
||||
|
||||
|
||||
@@ -43,6 +43,12 @@
|
||||
"type": "sqlite",
|
||||
"path": "/sqlite-data/workspace/data/app.db",
|
||||
"generated_id": "16678178-ff7e-4c97-8c83-0adeff214681"
|
||||
},
|
||||
{
|
||||
"name": "Test database 7 - SQLite DB",
|
||||
"type": "sqlite",
|
||||
"path": "/sqlite-data-2/workspace/data/app.db",
|
||||
"generated_id": "16678179-ff7e-4c97-8c83-0adeff214681"
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
+69
-68
@@ -12,12 +12,13 @@ services:
|
||||
- cargo-registry:/usr/local/cargo/registry
|
||||
- cargo-git:/usr/local/cargo/git
|
||||
# - cargo-target:/app/target
|
||||
- sqlite-data:/sqlite-data/workspace/data
|
||||
# - sqlite-data:/sqlite-data/workspace/data
|
||||
# - ./scripts/sqlite/test-db:/sqlite-data-2/workspace/data
|
||||
environment:
|
||||
APP_ENV: development
|
||||
LOG: debug
|
||||
TZ: "Europe/Paris"
|
||||
EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiZjg4Y2E0MDMtNDgwOS00NGM4LTlkZjItY2VkNWYwYzhkNTM2IiwibWFzdGVyS2V5QjY0IjoiQlhWM1hvbEM2NTZTVjdkTmdjV1BHUWxrKytycExJNmxHRGk3Q1BCNWllbz0ifQ=="
|
||||
EDGE_KEY: "eyJzZXJ2ZXJVcmwiOiJodHRwOi8vbG9jYWxob3N0Ojg4ODciLCJhZ2VudElkIjoiNjI1MDQzY2YtN2MwMC00M2M4LWJjYzktZDM1MTk5ODk2ZGNkIiwibWFzdGVyS2V5QjY0IjoiQlhWM1hvbEM2NTZTVjdkTmdjV1BHUWxrKytycExJNmxHRGk3Q1BCNWllbz0ifQ=="
|
||||
#POOLING: 1
|
||||
#DATABASES_CONFIG_FILE: "config.toml"
|
||||
extra_hosts:
|
||||
@@ -38,69 +39,69 @@ services:
|
||||
- POSTGRES_PASSWORD=changeme
|
||||
networks:
|
||||
- portabase
|
||||
#
|
||||
# db-mariadb:
|
||||
# container_name: db-mariadb
|
||||
# image: mariadb:latest
|
||||
# ports:
|
||||
# - "3311:3306"
|
||||
# environment:
|
||||
# - MYSQL_DATABASE=mariadb
|
||||
# - MYSQL_USER=mariadb
|
||||
# - MYSQL_PASSWORD=changeme
|
||||
# - MYSQL_RANDOM_ROOT_PASSWORD=yes
|
||||
# volumes:
|
||||
# - mariadb-data:/var/lib/mysql
|
||||
# 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
|
||||
|
||||
db-mariadb:
|
||||
container_name: db-mariadb
|
||||
image: mariadb:latest
|
||||
ports:
|
||||
- "3311:3306"
|
||||
environment:
|
||||
- MYSQL_DATABASE=mariadb
|
||||
- MYSQL_USER=mariadb
|
||||
- MYSQL_PASSWORD=changeme
|
||||
- MYSQL_RANDOM_ROOT_PASSWORD=yes
|
||||
volumes:
|
||||
- mariadb-data:/var/lib/mysql
|
||||
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
|
||||
|
||||
sqlite:
|
||||
container_name: db-sqlite
|
||||
image: keinos/sqlite3
|
||||
volumes:
|
||||
- sqlite-data:/workspace/data
|
||||
working_dir: /workspace
|
||||
command: tail -f /dev/null
|
||||
stdin_open: true
|
||||
tty: true
|
||||
# sqlite:
|
||||
# container_name: db-sqlite
|
||||
# image: keinos/sqlite3
|
||||
# volumes:
|
||||
# - sqlite-data:/workspace/data
|
||||
# working_dir: /workspace
|
||||
# command: tail -f /dev/null
|
||||
# stdin_open: true
|
||||
# tty: true
|
||||
|
||||
|
||||
volumes:
|
||||
@@ -109,10 +110,10 @@ volumes:
|
||||
# cargo-target:
|
||||
|
||||
postgres-data:
|
||||
mariadb-data:
|
||||
mongodb-data:
|
||||
mongodb-data-auth:
|
||||
sqlite-data:
|
||||
# mariadb-data:
|
||||
# mongodb-data:
|
||||
# mongodb-data-auth:
|
||||
# sqlite-data:
|
||||
|
||||
networks:
|
||||
portabase:
|
||||
|
||||
@@ -1,99 +0,0 @@
|
||||
#!/bin/bash
|
||||
|
||||
set -e
|
||||
|
||||
if [ -z "$1" ]; then
|
||||
echo "Usage: ./release <version>"
|
||||
echo "Example: ./release v1.0.0"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
VERSION=$1
|
||||
CURRENT_BRANCH=$(git rev-parse --abbrev-ref HEAD)
|
||||
|
||||
if [ "$CURRENT_BRANCH" = "main" ]; then
|
||||
if [[ ! "$VERSION" =~ ^[0-9]+\.[0-9]+\.[0-9]+$ ]]; then
|
||||
echo "Error: On 'main' branch, only release tags (X.Y.Z) are allowed."
|
||||
exit 1
|
||||
fi
|
||||
else
|
||||
if [[ ! "$VERSION" =~ ^[0-9]+\.[0-9]+\.[0-9]+-rc\.[0-9]+$ ]]; then
|
||||
echo "Error: On branch '$CURRENT_BRANCH', only RC tags matching X.Y.Z-rc.W are allowed (e.g., 1.0.0-rc.1)."
|
||||
exit 1
|
||||
fi
|
||||
fi
|
||||
|
||||
CLEAN_VERSION=${VERSION#v}
|
||||
CURRENT_DATE=$(date +%Y-%m-%d)
|
||||
|
||||
echo "Preparing release $VERSION..."
|
||||
|
||||
|
||||
# package.json
|
||||
if [ -f package.json ]; then
|
||||
echo "Updating package.json..."
|
||||
if sed --version >/dev/null 2>&1; then
|
||||
sed -i "s/\"version\": \".*\"/\"version\": \"$CLEAN_VERSION\"/" package.json
|
||||
else
|
||||
sed -i '' "s/\"version\": \".*\"/\"version\": \"$CLEAN_VERSION\"/" package.json
|
||||
fi
|
||||
fi
|
||||
|
||||
# pyproject.toml
|
||||
if [ -f pyproject.toml ]; then
|
||||
echo "Updating pyproject.toml..."
|
||||
if sed --version >/dev/null 2>&1; then
|
||||
sed -i "s/^version = \".*\"/version = \"$CLEAN_VERSION\"/" pyproject.toml
|
||||
else
|
||||
sed -i '' "s/^version = \".*\"/version = \"$CLEAN_VERSION\"/" pyproject.toml
|
||||
fi
|
||||
fi
|
||||
|
||||
# Cargo.toml
|
||||
if [ -f Cargo.toml ]; then
|
||||
echo "Updating Cargo.toml..."
|
||||
if sed --version >/dev/null 2>&1; then
|
||||
sed -i "s/^version = \".*\"/version = \"$CLEAN_VERSION\"/" Cargo.toml
|
||||
else
|
||||
sed -i '' "s/^version = \".*\"/version = \"$CLEAN_VERSION\"/" Cargo.toml
|
||||
fi
|
||||
fi
|
||||
|
||||
# CITATION.cff
|
||||
if [ -f CITATION.cff ]; then
|
||||
echo "Updating CITATION.cff..."
|
||||
if sed --version >/dev/null 2>&1; then
|
||||
sed -i "s/^version: .*/version: $CLEAN_VERSION/" CITATION.cff
|
||||
sed -i "s/^date-released: .*/date-released: \"$CURRENT_DATE\"/" CITATION.cff
|
||||
else
|
||||
sed -i '' "s/^version: .*/version: $CLEAN_VERSION/" CITATION.cff
|
||||
sed -i '' "s/^date-released: .*/date-released: \"$CURRENT_DATE\"/" CITATION.cff
|
||||
fi
|
||||
fi
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
git add .
|
||||
|
||||
if ! git diff-index --quiet HEAD --; then
|
||||
echo "Committing changes..."
|
||||
git commit -m "chore(release): $VERSION"
|
||||
else
|
||||
echo "No changes to commit. Proceeding to tag..."
|
||||
fi
|
||||
|
||||
if git rev-parse "$VERSION" >/dev/null 2>&1; then
|
||||
echo "Tag $VERSION already exists. Aborting."
|
||||
exit 1
|
||||
fi
|
||||
|
||||
echo "Creating tag $VERSION..."
|
||||
git tag -a "$VERSION" -m "Release $VERSION"
|
||||
|
||||
echo "Pushing changes and tags to remote..."
|
||||
git push
|
||||
git push origin "$VERSION"
|
||||
|
||||
echo "Successfully released $VERSION!"
|
||||
@@ -60,7 +60,7 @@ SELECT
|
||||
)
|
||||
FROM users u
|
||||
JOIN generate_series(1, 30) AS p(post_no)
|
||||
ON u.id <= 300000;
|
||||
ON u.id <= 500000;
|
||||
|
||||
-- ============================================================
|
||||
-- OPTIONAL: FORCE DISK MATERIALIZATION
|
||||
|
||||
Binary file not shown.
@@ -17,7 +17,7 @@ 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() {
|
||||
if cfg.username.is_empty() || cfg.password.is_empty() {
|
||||
Ok(format!("mongodb://{}:{}/{}", cfg.host, cfg.port, cfg.database))
|
||||
} else {
|
||||
Ok(format!(
|
||||
|
||||
@@ -1,2 +1,3 @@
|
||||
pub mod status;
|
||||
pub mod backup;
|
||||
pub mod backup;
|
||||
pub mod restore;
|
||||
@@ -0,0 +1,33 @@
|
||||
use crate::services::api::models::agent::restore::ResultRestoreResponse;
|
||||
use crate::services::api::{ApiClient, ApiError};
|
||||
use anyhow::Result;
|
||||
use reqwest::Method;
|
||||
use serde::Serialize;
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct ResultRestoreRequest {
|
||||
#[serde(rename = "generatedId")]
|
||||
pub generated_id: String,
|
||||
pub status: String,
|
||||
}
|
||||
|
||||
impl ApiClient {
|
||||
pub async fn restore_result(
|
||||
&self,
|
||||
agent_id: impl Into<String>,
|
||||
generated_id: impl Into<String>,
|
||||
status: impl Into<String>,
|
||||
) -> Result<Option<ResultRestoreResponse>, ApiError> {
|
||||
let body = ResultRestoreRequest {
|
||||
generated_id: generated_id.into(),
|
||||
status: status.into(),
|
||||
};
|
||||
|
||||
let agent_id = agent_id.into();
|
||||
|
||||
let path = format!("/agent/{}/restore", agent_id);
|
||||
|
||||
self.request_with_body(Method::POST, path.as_str(), &body)
|
||||
.await
|
||||
}
|
||||
}
|
||||
@@ -1,2 +1,3 @@
|
||||
pub mod status;
|
||||
pub mod backup;
|
||||
pub mod backup;
|
||||
pub mod restore;
|
||||
@@ -0,0 +1,7 @@
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
#[derive(Debug, Serialize, Deserialize)]
|
||||
pub struct ResultRestoreResponse {
|
||||
pub message: String,
|
||||
pub status: bool,
|
||||
}
|
||||
@@ -1,415 +0,0 @@
|
||||
#![allow(dead_code)]
|
||||
|
||||
use crate::core::context::Context as CoreContext;
|
||||
use crate::domain::factory::DatabaseFactory;
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::services::config::{DatabaseConfig, DatabasesConfig, DbType};
|
||||
use crate::services::storage;
|
||||
use crate::utils::common::BackupMethod;
|
||||
use crate::utils::compress::compress_to_tar_gz_large;
|
||||
use anyhow::Result;
|
||||
use futures::future::join_all;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
use tempfile::TempDir;
|
||||
use tracing::{error, info};
|
||||
use crate::utils::locks::FileLock;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct BackupResult {
|
||||
pub generated_id: String,
|
||||
pub db_type: DbType,
|
||||
pub status: String,
|
||||
pub backup_file: Option<PathBuf>,
|
||||
pub code: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct UploadResult {
|
||||
pub storage_id: String,
|
||||
pub success: bool,
|
||||
pub error: Option<String>,
|
||||
pub remote_file_path: Option<String>,
|
||||
pub total_size: Option<u64>,
|
||||
}
|
||||
|
||||
pub struct BackupService {
|
||||
ctx: Arc<CoreContext>,
|
||||
}
|
||||
|
||||
impl BackupService {
|
||||
pub fn new(ctx: Arc<CoreContext>) -> Self {
|
||||
Self { ctx }
|
||||
}
|
||||
|
||||
pub async fn dispatch(
|
||||
&self,
|
||||
generated_id: &String,
|
||||
config: &DatabasesConfig,
|
||||
method: BackupMethod,
|
||||
storages: &Vec<DatabaseStorage>,
|
||||
encrypt: bool,
|
||||
) {
|
||||
if let Some(cfg) = config
|
||||
.databases
|
||||
.iter()
|
||||
.find(|c| c.generated_id == generated_id.as_str())
|
||||
{
|
||||
let db_cfg = cfg.clone();
|
||||
let ctx = self.ctx.clone();
|
||||
let storages_clone = storages.clone();
|
||||
let generated_id_clone = generated_id.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
match TempDir::new() {
|
||||
Ok(temp_dir) => {
|
||||
match FileLock::is_locked(&generated_id_clone).await {
|
||||
Ok(true) => {
|
||||
error!("Backup already running for {}", &generated_id_clone);
|
||||
return;
|
||||
}
|
||||
Ok(false) => {
|
||||
match ctx
|
||||
.api
|
||||
.backup_create(
|
||||
method.clone().to_string(),
|
||||
ctx.edge_key.agent_id.clone(),
|
||||
&generated_id_clone,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(backup_created_result) => {
|
||||
info!("Backup created successfully");
|
||||
let tmp_path = temp_dir.path().to_path_buf();
|
||||
info!("Created temp directory {}", tmp_path.display());
|
||||
match BackupService::run(db_cfg, &tmp_path).await {
|
||||
Ok(mut result) => {
|
||||
let backup_id = backup_created_result.unwrap().backup.id;
|
||||
|
||||
if result.status == "failed" {
|
||||
error!("Backup failed early for {}", result.generated_id);
|
||||
let service = BackupService { ctx: ctx.clone() };
|
||||
let _ = service
|
||||
.send_result(result, vec![], &backup_id)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
if let Some(backup_file) = result.backup_file.take() {
|
||||
match compress_to_tar_gz_large(&backup_file).await {
|
||||
Ok(compression_result) => {
|
||||
result.backup_file =
|
||||
Some(compression_result.compressed_path);
|
||||
let service = BackupService { ctx: ctx.clone() };
|
||||
match service
|
||||
.upload(
|
||||
result.clone(),
|
||||
method,
|
||||
storages_clone.clone(),
|
||||
encrypt,
|
||||
&backup_id,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(upload_result) => {
|
||||
match service
|
||||
.send_result(
|
||||
result,
|
||||
upload_result,
|
||||
&backup_id,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => {
|
||||
return;
|
||||
}
|
||||
Err(e) => {
|
||||
error!(
|
||||
"Failed to send backup result: {}",
|
||||
e
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
error!(
|
||||
"Failed to upload backup files: {}",
|
||||
e
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
error!(
|
||||
"Failed to compress backup file : {}",
|
||||
e
|
||||
);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
error!("No backup file generated");
|
||||
}
|
||||
}
|
||||
Err(e) => error!("BackupService run failed: {}", e),
|
||||
}
|
||||
// TempDir is automatically deleted when dropped here
|
||||
}
|
||||
Err(e) => error!("Backup creation failed: {}", e),
|
||||
}
|
||||
}
|
||||
Err(e) => error!("An error occurred while checking lock : {}", e),
|
||||
}
|
||||
}
|
||||
Err(e) => error!("Failed to create temp dir: {}", e),
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn run(cfg: DatabaseConfig, tmp_path: &Path) -> Result<BackupResult> {
|
||||
let db_instance = DatabaseFactory::create_for_backup(cfg.clone()).await;
|
||||
let generated_id = cfg.generated_id.clone();
|
||||
let db_type = cfg.db_type.clone();
|
||||
|
||||
|
||||
let reachable = match db_instance.ping().await {
|
||||
Ok(v) => v,
|
||||
Err(e) => {
|
||||
error!("Ping failed: {}", e);
|
||||
return Err(e.into());
|
||||
}
|
||||
};
|
||||
|
||||
|
||||
info!("Reachable: {}", reachable);
|
||||
if !reachable {
|
||||
return Ok(BackupResult {
|
||||
generated_id,
|
||||
db_type,
|
||||
status: "failed".into(),
|
||||
backup_file: None,
|
||||
code: None,
|
||||
});
|
||||
}
|
||||
|
||||
match db_instance.backup(tmp_path).await {
|
||||
Ok(file) => Ok(BackupResult {
|
||||
generated_id,
|
||||
db_type,
|
||||
status: "success".into(),
|
||||
backup_file: Some(file),
|
||||
code: None,
|
||||
}),
|
||||
Err(e) => match e.to_string().as_str() {
|
||||
"backup_already_in_progress" => Ok(BackupResult {
|
||||
generated_id,
|
||||
db_type,
|
||||
status: "failed".into(),
|
||||
backup_file: None,
|
||||
code: Some(e.to_string()),
|
||||
}),
|
||||
_ => Ok(BackupResult {
|
||||
generated_id,
|
||||
db_type,
|
||||
status: "failed".into(),
|
||||
backup_file: None,
|
||||
code: None,
|
||||
}),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn upload(
|
||||
&self,
|
||||
result: BackupResult,
|
||||
method: BackupMethod,
|
||||
storages: Vec<DatabaseStorage>,
|
||||
encrypt: bool,
|
||||
backup_id: &String,
|
||||
) -> Result<Vec<UploadResult>> {
|
||||
if result.code.as_deref() == Some("backup_already_in_progress") {
|
||||
info!("Skipping send: backup already in progress");
|
||||
anyhow::bail!("backup_already_in_progres");
|
||||
}
|
||||
|
||||
let upload_futures = storages.into_iter().map(|storage| {
|
||||
info!(
|
||||
"Uploading storage -> {:?} for {:?}",
|
||||
storage.provider, storage.id
|
||||
);
|
||||
let provider = storage::get_provider(&storage);
|
||||
let result_clone = result.clone();
|
||||
let ctx_clone = self.ctx.clone();
|
||||
let storages_clone = storage.clone();
|
||||
let storage_id = storages_clone.id;
|
||||
let generated_id = result_clone.generated_id.clone();
|
||||
|
||||
async move {
|
||||
match self
|
||||
.ctx
|
||||
.api
|
||||
.backup_upload_init(
|
||||
self.ctx.edge_key.agent_id.clone(),
|
||||
generated_id.clone(),
|
||||
storage_id.clone(),
|
||||
backup_id,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(upload_init_result) => {
|
||||
info!("Uploading init result: {:#?}", upload_init_result);
|
||||
let backup_storage_id = upload_init_result.unwrap().backup_storage.id.clone();
|
||||
match provider {
|
||||
Some(provider) => {
|
||||
let upload_result = provider
|
||||
.upload(
|
||||
ctx_clone,
|
||||
result_clone,
|
||||
method,
|
||||
&storage,
|
||||
Some(encrypt),
|
||||
)
|
||||
.await;
|
||||
|
||||
let status = if upload_result.success {
|
||||
"success"
|
||||
} else {
|
||||
"failed"
|
||||
};
|
||||
|
||||
if status != "success" {
|
||||
return upload_result;
|
||||
}
|
||||
|
||||
info!("Storage {} uploaded to remote path {:?}", storage_id, upload_result.remote_file_path);
|
||||
|
||||
let (remote_path, total_size) = match (
|
||||
&upload_result.remote_file_path,
|
||||
upload_result.total_size,
|
||||
) {
|
||||
(Some(path), Some(size)) => (path.clone(), size),
|
||||
_ => {
|
||||
return UploadResult {
|
||||
storage_id: storage_id.clone(),
|
||||
success: false,
|
||||
error: Some("remote_file_path or total_size missing".to_string()),
|
||||
remote_file_path: None,
|
||||
total_size: None,
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
match self.ctx.api.backup_upload_status(
|
||||
self.ctx.edge_key.agent_id.clone(),
|
||||
generated_id.clone(),
|
||||
backup_storage_id,
|
||||
status,
|
||||
remote_path,
|
||||
total_size,
|
||||
backup_id,
|
||||
).await {
|
||||
Ok(_) => {
|
||||
upload_result
|
||||
}
|
||||
Err(err) => {
|
||||
error!(
|
||||
"backup_upload_status failed (generated_id={}, storage_id={}): {}",
|
||||
generated_id, storage_id, err
|
||||
);
|
||||
UploadResult {
|
||||
storage_id: storage_id.clone(),
|
||||
success: false,
|
||||
error: Some(err.to_string()),
|
||||
remote_file_path: None,
|
||||
total_size: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
None => {
|
||||
error!("Skipping storage due to missing provider");
|
||||
UploadResult {
|
||||
storage_id: storage_id.clone(),
|
||||
success: false,
|
||||
error: Some(
|
||||
"Skipping storage due to missing provider".to_string(),
|
||||
),
|
||||
remote_file_path: None,
|
||||
total_size: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
error!(
|
||||
"Unable to create the storage backup on remote server : {}",
|
||||
e
|
||||
);
|
||||
UploadResult {
|
||||
storage_id: storage_id.clone(),
|
||||
success: false,
|
||||
error: Some(
|
||||
"Unable to create the storage backup on remote server".to_string(),
|
||||
),
|
||||
remote_file_path: None,
|
||||
total_size: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
let results: Vec<UploadResult> = join_all(upload_futures).await;
|
||||
info!("Upload results: {:#?}", results);
|
||||
|
||||
Ok(results)
|
||||
}
|
||||
|
||||
pub async fn send_result(
|
||||
&self,
|
||||
result: BackupResult,
|
||||
upload_results: Vec<UploadResult>,
|
||||
backup_id: &String,
|
||||
) -> Result<()> {
|
||||
let status = if upload_results.iter().any(|r| r.success) {
|
||||
"success"
|
||||
} else {
|
||||
"failed"
|
||||
};
|
||||
|
||||
let file_size = if status == "failed" {
|
||||
None
|
||||
} else {
|
||||
let mut sum = 0u64;
|
||||
let mut count = 0u64;
|
||||
|
||||
for size in upload_results.iter().filter_map(|r| r.total_size) {
|
||||
sum += size;
|
||||
count += 1;
|
||||
}
|
||||
|
||||
if count == 0 {
|
||||
None
|
||||
} else {
|
||||
Some(sum / count)
|
||||
}
|
||||
};
|
||||
|
||||
match self
|
||||
.ctx
|
||||
.api
|
||||
.backup_update(self.ctx.edge_key.agent_id.clone(), backup_id, status, file_size, &result.generated_id)
|
||||
.await
|
||||
{
|
||||
Ok(_result) => Ok(()),
|
||||
Err(e) => {
|
||||
error!(
|
||||
"backup_update failed (generated_id={}, backup_id={}): {}",
|
||||
&result.generated_id, &backup_id, e
|
||||
);
|
||||
Err(e.into())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,20 @@
|
||||
use super::service::BackupService;
|
||||
use crate::utils::compress::compress_to_tar_gz_large;
|
||||
use std::path::PathBuf;
|
||||
use anyhow::Result;
|
||||
|
||||
impl BackupService {
|
||||
|
||||
pub async fn compress_backup(
|
||||
&self,
|
||||
backup_file: Option<PathBuf>,
|
||||
) -> Result<PathBuf> {
|
||||
|
||||
let file = backup_file
|
||||
.ok_or_else(|| anyhow::anyhow!("No backup file generated"))?;
|
||||
|
||||
let compression = compress_to_tar_gz_large(&file).await?;
|
||||
|
||||
Ok(compression.compressed_path)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,44 @@
|
||||
use super::service::BackupService;
|
||||
use crate::services::config::DatabasesConfig;
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::utils::common::BackupMethod;
|
||||
use tracing::error;
|
||||
|
||||
impl BackupService {
|
||||
|
||||
pub async fn dispatch(
|
||||
&self,
|
||||
generated_id: &String,
|
||||
config: &DatabasesConfig,
|
||||
method: BackupMethod,
|
||||
storages: &Vec<DatabaseStorage>,
|
||||
encrypt: bool,
|
||||
) {
|
||||
|
||||
let Some(cfg) = config
|
||||
.databases
|
||||
.iter()
|
||||
.find(|c| c.generated_id == generated_id.as_str())
|
||||
else {
|
||||
error!("Database config not found for {}", generated_id);
|
||||
return;
|
||||
};
|
||||
|
||||
let service = Self {
|
||||
ctx: self.ctx.clone(),
|
||||
};
|
||||
|
||||
let db_cfg = cfg.clone();
|
||||
let storages = storages.clone();
|
||||
let generated_id = generated_id.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = service
|
||||
.execute_backup(generated_id, db_cfg, method, storages, encrypt)
|
||||
.await
|
||||
{
|
||||
error!("Backup execution failed: {}", e);
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
use super::service::BackupService;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::utils::common::BackupMethod;
|
||||
use crate::utils::locks::FileLock;
|
||||
|
||||
use tempfile::TempDir;
|
||||
use anyhow::Result;
|
||||
|
||||
impl BackupService {
|
||||
|
||||
pub async fn execute_backup(
|
||||
&self,
|
||||
generated_id: String,
|
||||
db_cfg: DatabaseConfig,
|
||||
method: BackupMethod,
|
||||
storages: Vec<DatabaseStorage>,
|
||||
encrypt: bool,
|
||||
) -> Result<()> {
|
||||
|
||||
if FileLock::is_locked(&generated_id).await? {
|
||||
anyhow::bail!("backup already running");
|
||||
}
|
||||
|
||||
let backup = self.create_backup_record(&generated_id, &method).await?;
|
||||
let backup_id = backup.backup.id;
|
||||
|
||||
let temp_dir = TempDir::new()?;
|
||||
let tmp_path = temp_dir.path();
|
||||
|
||||
let mut result = Self::run(db_cfg, tmp_path).await?;
|
||||
|
||||
if result.status == "failed" {
|
||||
self.send_result(result, vec![], &backup_id).await?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let compressed = self.compress_backup(result.backup_file.take()).await?;
|
||||
result.backup_file = Some(compressed);
|
||||
|
||||
let uploads = self
|
||||
.upload(result.clone(), method, storages, encrypt, &backup_id)
|
||||
.await?;
|
||||
|
||||
self.send_result(result, uploads, &backup_id).await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
use super::service::BackupService;
|
||||
use crate::utils::common::BackupMethod;
|
||||
use anyhow::{Result, anyhow};
|
||||
use crate::services::api::models::agent::backup::BackupResponse;
|
||||
|
||||
impl BackupService {
|
||||
|
||||
pub async fn create_backup_record(
|
||||
&self,
|
||||
generated_id: &str,
|
||||
method: &BackupMethod,
|
||||
) -> Result<BackupResponse> {
|
||||
|
||||
let response = self
|
||||
.ctx
|
||||
.api
|
||||
.backup_create(
|
||||
method.to_string(),
|
||||
self.ctx.edge_key.agent_id.clone(),
|
||||
generated_id,
|
||||
)
|
||||
.await?;
|
||||
|
||||
response.ok_or_else(|| anyhow!("backup_create returned empty response"))
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
pub mod service;
|
||||
pub mod dispatcher;
|
||||
pub mod executor;
|
||||
pub mod compressor;
|
||||
pub mod uploader;
|
||||
pub mod result;
|
||||
pub mod models;
|
||||
pub mod helpers;
|
||||
pub mod runner;
|
||||
|
||||
pub use service::BackupService;
|
||||
@@ -0,0 +1,22 @@
|
||||
#![allow(dead_code)]
|
||||
|
||||
use std::path::PathBuf;
|
||||
use crate::services::config::DbType;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct BackupResult {
|
||||
pub generated_id: String,
|
||||
pub db_type: DbType,
|
||||
pub status: String,
|
||||
pub backup_file: Option<PathBuf>,
|
||||
pub code: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct UploadResult {
|
||||
pub storage_id: String,
|
||||
pub success: bool,
|
||||
pub error: Option<String>,
|
||||
pub remote_file_path: Option<String>,
|
||||
pub total_size: Option<u64>,
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
use super::models::{BackupResult, UploadResult};
|
||||
use super::service::BackupService;
|
||||
use crate::services::api::ApiError;
|
||||
use crate::services::api::models::agent::backup::BackupResponse;
|
||||
use anyhow::Result;
|
||||
use tracing::error;
|
||||
|
||||
impl BackupService {
|
||||
pub async fn send_result(
|
||||
&self,
|
||||
result: BackupResult,
|
||||
upload_results: Vec<UploadResult>,
|
||||
backup_id: &String,
|
||||
) -> Result<Option<BackupResponse>, ApiError> {
|
||||
let status = if upload_results.iter().any(|r| r.success) {
|
||||
"success"
|
||||
} else {
|
||||
"failed"
|
||||
};
|
||||
|
||||
let file_size = upload_results
|
||||
.iter()
|
||||
.filter_map(|r| r.total_size)
|
||||
.reduce(|a, b| a + b)
|
||||
.map(|sum| sum / upload_results.len() as u64);
|
||||
|
||||
self.ctx
|
||||
.api
|
||||
.backup_update(
|
||||
self.ctx.edge_key.agent_id.clone(),
|
||||
backup_id,
|
||||
status,
|
||||
file_size,
|
||||
&result.generated_id,
|
||||
)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
error!(
|
||||
"backup_update failed (generated_id={}, backup_id={}): {}",
|
||||
result.generated_id, backup_id, e
|
||||
);
|
||||
e.into()
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,70 @@
|
||||
use super::models::BackupResult;
|
||||
use super::service::BackupService;
|
||||
|
||||
use crate::domain::factory::DatabaseFactory;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
|
||||
use anyhow::Result;
|
||||
use std::path::Path;
|
||||
use tracing::{error, info};
|
||||
|
||||
impl BackupService {
|
||||
|
||||
pub async fn run(
|
||||
cfg: DatabaseConfig,
|
||||
tmp_path: &Path,
|
||||
) -> Result<BackupResult> {
|
||||
|
||||
let db = DatabaseFactory::create_for_backup(cfg.clone()).await;
|
||||
|
||||
let generated_id = cfg.generated_id.clone();
|
||||
let db_type = cfg.db_type.clone();
|
||||
|
||||
let reachable = match db.ping().await {
|
||||
Ok(v) => v,
|
||||
Err(e) => {
|
||||
error!("Ping failed: {}", e);
|
||||
return Err(e.into());
|
||||
}
|
||||
};
|
||||
|
||||
info!("Reachable: {}", reachable);
|
||||
|
||||
if !reachable {
|
||||
return Ok(BackupResult {
|
||||
generated_id,
|
||||
db_type,
|
||||
status: "failed".into(),
|
||||
backup_file: None,
|
||||
code: None,
|
||||
});
|
||||
}
|
||||
|
||||
match db.backup(tmp_path).await {
|
||||
|
||||
Ok(file) => Ok(BackupResult {
|
||||
generated_id,
|
||||
db_type,
|
||||
status: "success".into(),
|
||||
backup_file: Some(file),
|
||||
code: None,
|
||||
}),
|
||||
|
||||
Err(e) if e.to_string() == "backup_already_in_progress" => Ok(BackupResult {
|
||||
generated_id,
|
||||
db_type,
|
||||
status: "failed".into(),
|
||||
backup_file: None,
|
||||
code: Some("backup_already_in_progress".into()),
|
||||
}),
|
||||
|
||||
Err(_) => Ok(BackupResult {
|
||||
generated_id,
|
||||
db_type,
|
||||
status: "failed".into(),
|
||||
backup_file: None,
|
||||
code: None,
|
||||
}),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
use std::sync::Arc;
|
||||
use crate::core::context::Context as CoreContext;
|
||||
|
||||
pub struct BackupService {
|
||||
pub ctx: Arc<CoreContext>,
|
||||
}
|
||||
|
||||
impl BackupService {
|
||||
pub fn new(ctx: Arc<CoreContext>) -> Self {
|
||||
Self { ctx }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,181 @@
|
||||
use super::service::BackupService;
|
||||
use super::models::{BackupResult, UploadResult};
|
||||
|
||||
use crate::services::storage;
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::utils::common::BackupMethod;
|
||||
|
||||
use futures::future::join_all;
|
||||
use anyhow::{Result, bail};
|
||||
use tracing::{info, error};
|
||||
|
||||
impl BackupService {
|
||||
|
||||
pub async fn upload(
|
||||
&self,
|
||||
result: BackupResult,
|
||||
method: BackupMethod,
|
||||
storages: Vec<DatabaseStorage>,
|
||||
encrypt: bool,
|
||||
backup_id: &String,
|
||||
) -> Result<Vec<UploadResult>> {
|
||||
|
||||
if result.code.as_deref() == Some("backup_already_in_progress") {
|
||||
info!("Skipping send: backup already in progress");
|
||||
bail!("backup_already_in_progress");
|
||||
}
|
||||
|
||||
let ctx = self.ctx.clone();
|
||||
|
||||
let futures = storages.into_iter().map(|storage| {
|
||||
|
||||
let ctx_clone = ctx.clone();
|
||||
let result_clone = result.clone();
|
||||
let provider = storage::get_provider(&storage);
|
||||
|
||||
let storage_id = storage.id.clone();
|
||||
let generated_id = result_clone.generated_id.clone();
|
||||
|
||||
async move {
|
||||
|
||||
info!("Uploading storage -> {:?} for {:?}", storage.provider, storage_id);
|
||||
|
||||
/*
|
||||
INIT STEP
|
||||
*/
|
||||
let init = match ctx_clone.api
|
||||
.backup_upload_init(
|
||||
ctx_clone.edge_key.agent_id.clone(),
|
||||
generated_id.clone(),
|
||||
storage_id.clone(),
|
||||
backup_id,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(v) => v,
|
||||
Err(e) => {
|
||||
error!("backup_upload_init failed: {}", e);
|
||||
|
||||
return UploadResult {
|
||||
storage_id,
|
||||
success: false,
|
||||
error: Some("backup_upload_init failed".into()),
|
||||
remote_file_path: None,
|
||||
total_size: None,
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
let backup_storage_id = match init {
|
||||
Some(v) => v.backup_storage.id,
|
||||
None => {
|
||||
return UploadResult {
|
||||
storage_id,
|
||||
success: false,
|
||||
error: Some("backup_upload_init returned empty response".into()),
|
||||
remote_file_path: None,
|
||||
total_size: None,
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
/*
|
||||
PROVIDER CHECK
|
||||
*/
|
||||
let Some(provider) = provider else {
|
||||
error!("Skipping storage due to missing provider");
|
||||
|
||||
return UploadResult {
|
||||
storage_id,
|
||||
success: false,
|
||||
error: Some("missing provider".into()),
|
||||
remote_file_path: None,
|
||||
total_size: None,
|
||||
};
|
||||
};
|
||||
|
||||
/*
|
||||
STORAGE UPLOAD
|
||||
*/
|
||||
let upload_result = provider
|
||||
.upload(
|
||||
ctx_clone.clone(),
|
||||
result_clone,
|
||||
method,
|
||||
&storage,
|
||||
Some(encrypt),
|
||||
)
|
||||
.await;
|
||||
|
||||
let status = if upload_result.success { "success" } else { "failed" };
|
||||
|
||||
if status != "success" {
|
||||
return upload_result;
|
||||
}
|
||||
|
||||
info!(
|
||||
"Storage {} uploaded to remote path {:?}",
|
||||
storage_id,
|
||||
upload_result.remote_file_path
|
||||
);
|
||||
|
||||
/*
|
||||
METADATA VALIDATION
|
||||
*/
|
||||
let (remote_path, total_size) = match (
|
||||
&upload_result.remote_file_path,
|
||||
upload_result.total_size,
|
||||
) {
|
||||
(Some(path), Some(size)) => (path.clone(), size),
|
||||
_ => {
|
||||
return UploadResult {
|
||||
storage_id,
|
||||
success: false,
|
||||
error: Some("remote_file_path or total_size missing".into()),
|
||||
remote_file_path: None,
|
||||
total_size: None,
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
/*
|
||||
STATUS UPDATE
|
||||
*/
|
||||
match ctx_clone.api.backup_upload_status(
|
||||
ctx_clone.edge_key.agent_id.clone(),
|
||||
generated_id,
|
||||
backup_storage_id,
|
||||
status,
|
||||
remote_path,
|
||||
total_size,
|
||||
backup_id,
|
||||
).await {
|
||||
|
||||
Ok(_) => upload_result,
|
||||
|
||||
Err(err) => {
|
||||
error!(
|
||||
"backup_upload_status failed (storage_id={}): {}",
|
||||
storage_id,
|
||||
err
|
||||
);
|
||||
|
||||
UploadResult {
|
||||
storage_id,
|
||||
success: false,
|
||||
error: Some(err.to_string()),
|
||||
remote_file_path: None,
|
||||
total_size: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
let results: Vec<UploadResult> = join_all(futures).await;
|
||||
|
||||
info!("Upload results: {:#?}", results);
|
||||
|
||||
Ok(results)
|
||||
}
|
||||
}
|
||||
@@ -1,246 +0,0 @@
|
||||
#![allow(dead_code)]
|
||||
#![warn(unused_assignments)]
|
||||
use crate::core::context::Context;
|
||||
use crate::domain::factory::DatabaseFactory;
|
||||
use crate::services::api::models::agent::status::DatabaseStatus;
|
||||
use crate::services::config::{DatabaseConfig, DatabasesConfig};
|
||||
use crate::utils::compress::decompress_large_tar_gz;
|
||||
use crate::utils::file::decrypt_file_stream_gcm;
|
||||
use anyhow::Result;
|
||||
use reqwest::{Client, Url};
|
||||
use serde::Serialize;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Arc;
|
||||
use tempfile::TempDir;
|
||||
use tracing::{error, info};
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
pub struct RestoreResult {
|
||||
#[serde(rename = "generatedId")]
|
||||
pub generated_id: String,
|
||||
pub status: String,
|
||||
}
|
||||
|
||||
pub struct RestoreService {
|
||||
ctx: Arc<Context>,
|
||||
}
|
||||
|
||||
impl RestoreService {
|
||||
pub fn new(ctx: Arc<Context>) -> Self {
|
||||
Self { ctx }
|
||||
}
|
||||
|
||||
pub async fn dispatch(&self, db: &DatabaseStatus, config: &DatabasesConfig) {
|
||||
if let Some(cfg) = config
|
||||
.databases
|
||||
.iter()
|
||||
.find(|c| c.generated_id == db.generated_id)
|
||||
{
|
||||
let db_cfg = cfg.clone();
|
||||
let ctx_clone = self.ctx.clone();
|
||||
let file_to_restore = db.data.restore.file.clone();
|
||||
if file_to_restore.is_none() {
|
||||
error!("restore file not found");
|
||||
return;
|
||||
}
|
||||
tokio::spawn(async move {
|
||||
match TempDir::new() {
|
||||
Ok(temp_dir) => {
|
||||
let tmp_path = temp_dir.path().to_path_buf();
|
||||
info!("Created temp directory {}", tmp_path.display());
|
||||
match RestoreService::run(
|
||||
&ctx_clone,
|
||||
db_cfg,
|
||||
&tmp_path,
|
||||
&file_to_restore.unwrap(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(result) => {
|
||||
let service = RestoreService { ctx: ctx_clone };
|
||||
service.send_result(result).await;
|
||||
}
|
||||
Err(e) => error!("Restoration error {}", e),
|
||||
}
|
||||
// TempDir is automatically deleted when dropped
|
||||
}
|
||||
Err(e) => error!("Failed to create temp dir: {}", e),
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn run(
|
||||
ctx: &Arc<Context>,
|
||||
cfg: DatabaseConfig,
|
||||
tmp_path: &Path,
|
||||
file_url: &str,
|
||||
) -> Result<RestoreResult> {
|
||||
let generated_id = cfg.generated_id.clone();
|
||||
|
||||
info!("File url: {}", file_url);
|
||||
|
||||
let client = Client::new();
|
||||
let response = client.get(file_url).send().await?;
|
||||
|
||||
if !response.status().is_success() {
|
||||
error!("Backup download failed with status {}", response.status());
|
||||
return Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "failed".into(),
|
||||
});
|
||||
}
|
||||
|
||||
let filename_from_header = response
|
||||
.headers()
|
||||
.get(reqwest::header::CONTENT_DISPOSITION)
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.and_then(|s| s.split("filename=").nth(1))
|
||||
.map(|f| f.trim_matches('"').to_string());
|
||||
|
||||
let filename_from_url = Url::parse(file_url).ok().and_then(|u| {
|
||||
u.path_segments()?
|
||||
.last()
|
||||
.filter(|s| !s.is_empty())
|
||||
.map(|s| s.to_string())
|
||||
});
|
||||
|
||||
let filename = filename_from_header
|
||||
.or(filename_from_url)
|
||||
.unwrap_or_else(|| "downloaded_file".to_string());
|
||||
|
||||
info!("File name: {}", filename);
|
||||
|
||||
let bytes = response.bytes().await?;
|
||||
|
||||
let is_legacy_file = if filename.ends_with(".sql") {
|
||||
true
|
||||
} else if filename.ends_with(".dump") {
|
||||
true
|
||||
} else {
|
||||
false
|
||||
};
|
||||
let downloaded_file = tmp_path.join(&filename);
|
||||
tokio::fs::write(&downloaded_file, &bytes).await?;
|
||||
info!("Backup downloaded to {}", downloaded_file.display());
|
||||
|
||||
let backup_file_path: PathBuf = if !is_legacy_file {
|
||||
let encrypted = if filename.ends_with(".tar.gz") {
|
||||
false
|
||||
} else if filename.ends_with(".tar.gz.enc") {
|
||||
true
|
||||
} else {
|
||||
return Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "failed".into(),
|
||||
});
|
||||
};
|
||||
|
||||
info!("Encrypted: {}", encrypted);
|
||||
|
||||
let mut compressed_archive = downloaded_file.clone();
|
||||
|
||||
if encrypted {
|
||||
let new_name = downloaded_file
|
||||
.file_name()
|
||||
.and_then(|n| n.to_str())
|
||||
.and_then(|n| n.strip_suffix(".enc"))
|
||||
.ok_or_else(|| anyhow::anyhow!("Invalid encrypted filename"))?;
|
||||
|
||||
let new_compressed_archive = tmp_path.join(new_name);
|
||||
|
||||
decrypt_file_stream_gcm(
|
||||
downloaded_file,
|
||||
new_compressed_archive.clone(),
|
||||
ctx.edge_key.master_key_b64.clone(),
|
||||
)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
error!("Failed to decrypt file: {}", e);
|
||||
e
|
||||
})?;
|
||||
|
||||
compressed_archive = new_compressed_archive;
|
||||
}
|
||||
|
||||
let decompressed_files =
|
||||
decompress_large_tar_gz(compressed_archive.as_path(), tmp_path).await?;
|
||||
|
||||
if decompressed_files.is_empty() {
|
||||
return Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "failed".into(),
|
||||
});
|
||||
}
|
||||
|
||||
if decompressed_files.len() == 1 {
|
||||
decompressed_files[0].clone()
|
||||
} else {
|
||||
compressed_archive
|
||||
}
|
||||
} else {
|
||||
downloaded_file.clone()
|
||||
};
|
||||
|
||||
let db_instance = DatabaseFactory::create_for_restore(cfg.clone(), &backup_file_path).await;
|
||||
let reachable = db_instance.ping().await.unwrap_or(false);
|
||||
info!("Reachable: {}", reachable);
|
||||
if !reachable {
|
||||
return Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "failed".into(),
|
||||
});
|
||||
}
|
||||
|
||||
match db_instance.restore(&backup_file_path).await {
|
||||
Ok(_) => Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "success".into(),
|
||||
}),
|
||||
Err(e) => {
|
||||
error!("Restore failed: {:?}", e);
|
||||
Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "failed".into(),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TODO : update with ctx api manager
|
||||
pub async fn send_result(&self, result: RestoreResult) {
|
||||
info!(
|
||||
"[RestoreService] DB: {} | Status: {}",
|
||||
result.generated_id, result.status,
|
||||
);
|
||||
|
||||
let client = reqwest::Client::new();
|
||||
let url = format!(
|
||||
"{}/api/agent/{}/restore",
|
||||
self.ctx.edge_key.server_url, self.ctx.edge_key.agent_id
|
||||
);
|
||||
|
||||
let body = RestoreResult {
|
||||
generated_id: result.generated_id,
|
||||
status: result.status,
|
||||
};
|
||||
|
||||
match client.post(&url).json(&body).send().await {
|
||||
Ok(resp) => {
|
||||
let status = resp.status();
|
||||
if status.is_success() {
|
||||
info!("Restoration result sent successfully");
|
||||
} else {
|
||||
let text = resp.text().await.unwrap_or_default();
|
||||
error!(
|
||||
"Restoration result failed, status: {}, body: {}",
|
||||
status, text
|
||||
);
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Failed to send restoration result: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
use super::service::RestoreService;
|
||||
|
||||
use crate::utils::compress::decompress_large_tar_gz;
|
||||
use crate::utils::file::decrypt_file_stream_gcm;
|
||||
|
||||
use anyhow::Result;
|
||||
use std::path::{Path, PathBuf};
|
||||
|
||||
impl RestoreService {
|
||||
|
||||
pub async fn prepare_archive(
|
||||
&self,
|
||||
downloaded_file: PathBuf,
|
||||
tmp_path: &Path,
|
||||
) -> Result<PathBuf> {
|
||||
|
||||
let filename = downloaded_file
|
||||
.file_name()
|
||||
.unwrap()
|
||||
.to_string_lossy()
|
||||
.to_string();
|
||||
|
||||
let is_legacy = filename.ends_with(".sql") || filename.ends_with(".dump");
|
||||
|
||||
if is_legacy {
|
||||
return Ok(downloaded_file);
|
||||
}
|
||||
|
||||
let encrypted = filename.ends_with(".tar.gz.enc");
|
||||
|
||||
let mut archive = downloaded_file.clone();
|
||||
|
||||
if encrypted {
|
||||
|
||||
let new_name = filename.strip_suffix(".enc").unwrap();
|
||||
|
||||
let decrypted = tmp_path.join(new_name);
|
||||
|
||||
decrypt_file_stream_gcm(
|
||||
downloaded_file,
|
||||
decrypted.clone(),
|
||||
self.ctx.edge_key.master_key_b64.clone(),
|
||||
)
|
||||
.await?;
|
||||
|
||||
archive = decrypted;
|
||||
}
|
||||
|
||||
let files = decompress_large_tar_gz(archive.as_path(), tmp_path).await?;
|
||||
|
||||
if files.is_empty() {
|
||||
anyhow::bail!("archive empty");
|
||||
}
|
||||
|
||||
if files.len() == 1 {
|
||||
Ok(files[0].clone())
|
||||
} else {
|
||||
Ok(archive)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,42 @@
|
||||
use super::service::RestoreService;
|
||||
use crate::services::config::DatabasesConfig;
|
||||
use crate::services::api::models::agent::status::DatabaseStatus;
|
||||
|
||||
use tracing::error;
|
||||
|
||||
impl RestoreService {
|
||||
|
||||
pub async fn dispatch(&self, db: &DatabaseStatus, config: &DatabasesConfig) {
|
||||
|
||||
let Some(cfg) = config
|
||||
.databases
|
||||
.iter()
|
||||
.find(|c| c.generated_id == db.generated_id)
|
||||
else {
|
||||
error!("Database config not found");
|
||||
return;
|
||||
};
|
||||
|
||||
let Some(file_to_restore) = db.data.restore.file.clone() else {
|
||||
error!("restore file not found");
|
||||
return;
|
||||
};
|
||||
|
||||
let service = Self {
|
||||
ctx: self.ctx.clone(),
|
||||
};
|
||||
|
||||
let db_cfg = cfg.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
|
||||
if let Err(e) = service
|
||||
.execute_restore(db_cfg, file_to_restore)
|
||||
.await
|
||||
{
|
||||
error!("Restore failed: {}", e);
|
||||
}
|
||||
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,53 @@
|
||||
use super::service::RestoreService;
|
||||
|
||||
use reqwest::{Client, Url};
|
||||
use anyhow::Result;
|
||||
use std::path::{Path, PathBuf};
|
||||
use tracing::info;
|
||||
|
||||
impl RestoreService {
|
||||
|
||||
pub async fn download_backup(
|
||||
&self,
|
||||
file_url: &str,
|
||||
tmp_path: &Path,
|
||||
) -> Result<PathBuf> {
|
||||
|
||||
let client = Client::new();
|
||||
|
||||
let response = client.get(file_url).send().await?;
|
||||
|
||||
if !response.status().is_success() {
|
||||
anyhow::bail!("download failed");
|
||||
}
|
||||
|
||||
let filename_from_header = response
|
||||
.headers()
|
||||
.get(reqwest::header::CONTENT_DISPOSITION)
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.and_then(|s| s.split("filename=").nth(1))
|
||||
.map(|f| f.trim_matches('"').to_string());
|
||||
|
||||
|
||||
let filename_from_url = Url::parse(file_url).ok().and_then(|u| {
|
||||
u.path_segments()?
|
||||
.last()
|
||||
.filter(|s| !s.is_empty())
|
||||
.map(|s| s.to_string())
|
||||
});
|
||||
|
||||
let filename = filename_from_header
|
||||
.or(filename_from_url)
|
||||
.unwrap_or_else(|| "downloaded_file".to_string());
|
||||
|
||||
let path = tmp_path.join(&filename);
|
||||
|
||||
let bytes = response.bytes().await?;
|
||||
|
||||
tokio::fs::write(&path, &bytes).await?;
|
||||
|
||||
info!("Backup downloaded to {}", path.display());
|
||||
|
||||
Ok(path)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
use super::service::RestoreService;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
use tempfile::TempDir;
|
||||
use anyhow::Result;
|
||||
use tracing::info;
|
||||
|
||||
impl RestoreService {
|
||||
|
||||
pub async fn execute_restore(
|
||||
&self,
|
||||
cfg: DatabaseConfig,
|
||||
file_url: String,
|
||||
) -> Result<()> {
|
||||
|
||||
let temp_dir = TempDir::new()?;
|
||||
let tmp_path = temp_dir.path();
|
||||
|
||||
info!("Created temp directory {}", tmp_path.display());
|
||||
|
||||
let downloaded = self.download_backup(&file_url, tmp_path).await?;
|
||||
|
||||
let backup_file = self.prepare_archive(downloaded, tmp_path).await?;
|
||||
|
||||
let result = self.run_restore(cfg, backup_file).await?;
|
||||
|
||||
self.send_result(result).await;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,10 @@
|
||||
pub mod service;
|
||||
pub mod dispatcher;
|
||||
pub mod executor;
|
||||
pub mod downloader;
|
||||
pub mod archive;
|
||||
pub mod runner;
|
||||
pub mod result;
|
||||
pub mod models;
|
||||
|
||||
pub use service::RestoreService;
|
||||
@@ -0,0 +1,8 @@
|
||||
use serde::Serialize;
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
pub struct RestoreResult {
|
||||
#[serde(rename = "generatedId")]
|
||||
pub generated_id: String,
|
||||
pub status: String,
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
use super::service::RestoreService;
|
||||
use super::models::RestoreResult;
|
||||
|
||||
use tracing::{info, error};
|
||||
|
||||
impl RestoreService {
|
||||
pub async fn send_result(&self, result: RestoreResult) {
|
||||
|
||||
info!(
|
||||
"[RestoreService] DB: {} | Status: {}",
|
||||
result.generated_id, result.status
|
||||
);
|
||||
|
||||
match self.ctx
|
||||
.api
|
||||
.restore_result(
|
||||
self.ctx.edge_key.agent_id.clone(),
|
||||
&result.generated_id,
|
||||
&result.status,
|
||||
)
|
||||
.await {
|
||||
Ok(_) => {
|
||||
info!("Restoration result sent successfully");
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Failed to send restoration result: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,51 @@
|
||||
use super::service::RestoreService;
|
||||
use super::models::RestoreResult;
|
||||
|
||||
use crate::domain::factory::DatabaseFactory;
|
||||
use crate::services::config::DatabaseConfig;
|
||||
|
||||
use anyhow::Result;
|
||||
use std::path::PathBuf;
|
||||
use tracing::{info, error};
|
||||
|
||||
impl RestoreService {
|
||||
|
||||
pub async fn run_restore(
|
||||
&self,
|
||||
cfg: DatabaseConfig,
|
||||
backup_file: PathBuf,
|
||||
) -> Result<RestoreResult> {
|
||||
|
||||
let generated_id = cfg.generated_id.clone();
|
||||
|
||||
let db = DatabaseFactory::create_for_restore(cfg.clone(), &backup_file).await;
|
||||
|
||||
let reachable = db.ping().await.unwrap_or(false);
|
||||
|
||||
info!("Reachable: {}", reachable);
|
||||
|
||||
if !reachable {
|
||||
return Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "failed".into(),
|
||||
});
|
||||
}
|
||||
|
||||
match db.restore(&backup_file).await {
|
||||
|
||||
Ok(_) => Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "success".into(),
|
||||
}),
|
||||
|
||||
Err(e) => {
|
||||
error!("Restore failed: {:?}", e);
|
||||
|
||||
Ok(RestoreResult {
|
||||
generated_id,
|
||||
status: "failed".into(),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
use std::sync::Arc;
|
||||
use crate::core::context::Context;
|
||||
|
||||
pub struct RestoreService {
|
||||
pub ctx: Arc<Context>,
|
||||
}
|
||||
|
||||
impl RestoreService {
|
||||
pub fn new(ctx: Arc<Context>) -> Self {
|
||||
Self { ctx }
|
||||
}
|
||||
}
|
||||
@@ -1,7 +1,6 @@
|
||||
pub mod providers;
|
||||
|
||||
use crate::core::context::Context;
|
||||
use crate::services::backup::{BackupResult, UploadResult};
|
||||
use crate::utils::common::BackupMethod;
|
||||
use async_trait::async_trait;
|
||||
use providers::local;
|
||||
@@ -10,6 +9,7 @@ use providers::google_drive;
|
||||
use std::sync::Arc;
|
||||
use tracing::{error, info};
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::services::backup::models::{BackupResult, UploadResult};
|
||||
|
||||
#[async_trait]
|
||||
pub trait StorageProvider: Send + Sync {
|
||||
|
||||
@@ -3,7 +3,6 @@ mod models;
|
||||
|
||||
use crate::core::context::Context;
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::services::backup::{BackupResult, UploadResult};
|
||||
use crate::services::storage::StorageProvider;
|
||||
use crate::utils::common::BackupMethod;
|
||||
use crate::utils::file::{full_file_name, full_file_path};
|
||||
@@ -12,6 +11,7 @@ use async_trait::async_trait;
|
||||
use std::sync::Arc;
|
||||
use tokio::fs;
|
||||
use tracing::{error, info};
|
||||
use crate::services::backup::models::{BackupResult, UploadResult};
|
||||
use crate::services::storage::providers::google_drive::helpers::{upload_stream_to_google_drive};
|
||||
use crate::services::storage::providers::google_drive::models::GoogleDriveProviderConfig;
|
||||
|
||||
|
||||
@@ -1,6 +1,5 @@
|
||||
use crate::core::context::Context;
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::services::backup::{BackupResult, UploadResult};
|
||||
use crate::services::storage::StorageProvider;
|
||||
use crate::utils::common::BackupMethod;
|
||||
use crate::utils::file::{full_file_name, full_file_path};
|
||||
@@ -11,6 +10,7 @@ use reqwest::header::{HeaderMap, HeaderValue};
|
||||
use std::sync::Arc;
|
||||
use tokio::fs;
|
||||
use tracing::error;
|
||||
use crate::services::backup::models::{BackupResult, UploadResult};
|
||||
|
||||
pub struct LocalProvider;
|
||||
|
||||
|
||||
@@ -2,7 +2,6 @@ mod models;
|
||||
|
||||
use crate::core::context::Context;
|
||||
use crate::services::api::models::agent::status::DatabaseStorage;
|
||||
use crate::services::backup::{BackupResult, UploadResult};
|
||||
use crate::services::storage::StorageProvider;
|
||||
use crate::services::storage::providers::s3::models::S3ProviderConfig;
|
||||
use crate::utils::common::BackupMethod;
|
||||
@@ -19,6 +18,7 @@ use std::pin::Pin;
|
||||
use std::sync::Arc;
|
||||
use tokio::fs;
|
||||
use tracing::{error, info};
|
||||
use crate::services::backup::models::{BackupResult, UploadResult};
|
||||
|
||||
pub struct S3Provider {}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user