17 Commits
Author SHA1 Message Date
jpmschweitzerandClaude 0865763597 chore: release v1.4.0
Build and Push / release (push) Successful in 3s
Build and Push / build (push) Successful in 32s
Ships the Portainer backup executor.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-08 19:56:44 +02:00
jpmschweitzerandClaude 933196c2a2 feat(executors): back up Portainer's own state
Portainer keeps every stack definition, endpoint, user and access-control
rule in a BoltDB inside the portainer_data Docker volume. That volume
sits under /var/lib/docker/volumes/, and the daily config backup covers
~/docker-data and code-server-config only — so the thing that defines all
24 stacks was the one thing not backed up.

Calls Portainer's /api/backup rather than tarring the volume. BoltDB is a
single memory-mapped file, so copying it while Portainer writes can
capture a torn page; the API serialises a consistent snapshot.

A 200 whose body is not a readable archive is treated as failure. An
archive that will not open is worse than a missing one, because it looks
like a backup until the day it is needed. Writing that check found a real
gap in it: a truncated tar.gz raises EOFError, which is neither TarError
nor OSError, so the first version of the guard let it through.

Archives contain TLS certificates and private keys and are written 0600.
Retention only ever deletes files matching the exact name this executor
writes, so an unrelated archive left in the same directory survives.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-08 19:56:44 +02:00
jpmschweitzerandClaude 23cd5ddca8 chore: release v1.3.0
Build and Push / release (push) Successful in 4s
Build and Push / build (push) Successful in 2m7s
Carries the two new pruning executors and the POST /tasks fix, which has
been on main unreleased since the image only rebuilds on a version tag.

Also backfills the missing 1.2.0 changelog entry: that version was tagged
and shipped without one.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-08 11:01:52 +02:00
jpmschweitzerandClaude cfabe1b4d1 docs: correct the executor list in TASK_REGISTRATION
The "Other Executors" section advertised shell, python and docker
executors that were never implemented, and omitted every executor that
was. The missing shell executor in particular sent a recent piece of
work down the wrong path before the gap was noticed.

Lists the modules that actually exist and documents the config for the
two new ones.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-08 11:00:42 +02:00
jpmschweitzerandClaude 7b80e30691 feat(executors): add docker prune executor
Non-interactive equivalent of system-admin-toj's prune-docker.sh, which
prompts per stage and so cannot run from cron.

Only the stages that discard regenerable data run by default: build
cache and dangling images. Unused images and volumes are opt-in, because
docker volume prune removes volumes belonging to merely-stopped
containers rather than only orphaned ones, which on this host is a
plausible way to lose a database.

A failing stage is reported and the remaining stages still run, since a
partial reclaim beats none, but the task still ends up failed so the
error is not swallowed.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-08 11:00:42 +02:00
jpmschweitzerandClaude 788c03514a feat(executors): add postgres retention executor
Deletes rows past a retention window from a table on the shared Postgres
server. Written for sysmon's check_history, which grows with every
monitoring check and had no retention at all despite the docs promising
a 30-day rolling window.

Connects with the Scheduler's own credentials and overrides only the
database name, so no second set of secrets enters the stack. The target
database grants scheduler_user just SELECT and DELETE on the table, so a
bug here can drop old rows but cannot corrupt or forge history.

Table and column names cannot be bound as query parameters, so both are
validated against a strict identifier pattern before interpolation, and
a retention window below 1 day is refused rather than silently emptying
the table.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-08 11:00:42 +02:00
jpmschweitzerandClaude ae4f9e6a20 fix(api): return created task instead of 500 on POST /tasks
TaskResponse declared created_at and updated_at as str, but both are
timestamp columns and psycopg2 returns datetime objects. Pydantic
rejected every response, so the endpoint raised ResponseValidationError
after the INSERT had already committed.

Every task creation therefore looked like a failure, and the natural
retry failed again with a genuine duplicate-key violation, making it
appear the first attempt had done nothing.

Declaring them as datetime leaves the JSON on the wire unchanged
(FastAPI serialises to ISO 8601) and matches what GET /tasks/{task_name}
already returned.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-07 17:14:16 +02:00
jpmschweitzerandClaude Fable 5 c1fbc1cdb0 chore(ci): push images via git.schweitz.net registry
The .internal registry domain is being retired; git.schweitz.net now
serves the registry without SSO on /v2/.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-19 11:10:08 +02:00
jpmschweitzerandClaude Opus 4.6 22ec600ef3 feat(backup): add GCS offsite backup executor
Build and Push / release (push) Successful in 28s
Build and Push / build (push) Successful in 6m56s
Add gcs_backup_executor with git_bundle mode for backing up bare git
repos to Google Cloud Storage. Includes retention management and
bundle verification. Adds google-cloud-storage dependency and
GCS_CREDENTIALS_FILE setting.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-03-30 21:16:45 +02:00
jpmschweitzerandClaude Opus 4.5 31802a1281 chore: Bump version to 1.1.3
Build and Push / release (push) Successful in 3s
Build and Push / build (push) Successful in 5m45s
Test release to validate CI/CD auto-deploy workflow

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2026-01-08 12:06:53 +01:00
jpmschweitzerandClaude Opus 4.5 ec200c66ba fix(ci): Use curl for release creation
Build and Push / release (push) Successful in 3s
Build and Push / build (push) Successful in 1m10s
The release-action requires Go which isn't in the runner image.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2026-01-03 20:33:47 +01:00
jpmschweitzerandClaude Opus 4.5 2a6a91ed67 chore: Bump version to 1.1.1
Build and Push / release (push) Failing after 7s
Build and Push / build (push) Has been skipped
🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2026-01-03 20:30:12 +01:00
jpmschweitzerandClaude Opus 4.5 f334753918 ci: Auto-create release on version tag push
Change workflow trigger from manual release to tag push (v*).
Adds release job that creates Gitea release before building.

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2026-01-03 20:28:55 +01:00
jpmschweitzer 8a6bdb9547 update to AGENTS.md 2025-12-25 10:05:58 +01:00
jpmschweitzerandClaude Opus 4.5 424c94d79a feat: Add Gitea release cleanup executor
Build and Push / build (release) Successful in 27s
- Add gitea_release_cleanup_executor for automated release cleanup
- Add GITEA_TOKEN setting for API token authentication
- Configurable retention count, repo exclusions, and dry-run mode
- Designed to run daily before Watchtower (3 AM)

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-14 13:25:02 +01:00
jpmschweitzerandClaude Opus 4.5 8f1a4402e7 fix(ci): Correct Watchtower port
Build and Push / build (release) Successful in 28s
🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-14 12:57:08 +01:00
jpmschweitzerandClaude Opus 4.5 9ae10af1df ci: Add Watchtower update trigger after build
🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
2025-12-14 12:45:06 +01:00
16 changed files with 1496 additions and 13 deletions
+24 -5
View File
@@ -1,19 +1,32 @@
name: Build and Push
on:
release:
types: [published]
push:
tags:
- 'v*'
jobs:
release:
runs-on: ubuntu-latest
steps:
- name: Create Gitea Release
run: |
curl -sf -X POST \
-H "Authorization: token ${{ secrets.GITHUB_TOKEN }}" \
-H "Content-Type: application/json" \
-d '{"tag_name": "${{ github.ref_name }}", "name": "Release ${{ github.ref_name }}", "body": "Automated release for ${{ github.ref_name }}"}' \
"${{ github.server_url }}/api/v1/repos/${{ github.repository }}/releases"
build:
runs-on: ubuntu-latest
needs: release
steps:
- uses: actions/checkout@v4
- name: Login to Gitea Registry
uses: docker/login-action@v3
with:
registry: git.schweitz.internal
registry: git.schweitz.net
username: ${{ secrets.REGISTRY_USER }}
password: ${{ secrets.REGISTRY_PASSWORD }}
@@ -23,5 +36,11 @@ jobs:
context: .
push: true
tags: |
git.schweitz.internal/jpmschweitzer/scheduler:latest
git.schweitz.internal/jpmschweitzer/scheduler:${{ github.ref_name }}
git.schweitz.net/jpmschweitzer/scheduler:latest
git.schweitz.net/jpmschweitzer/scheduler:${{ github.ref_name }}
- name: Trigger Watchtower update
if: success()
run: |
curl -sf -H "Authorization: Bearer ${{ secrets.WATCHTOWER_TOKEN }}" \
http://watchtower:8080/v1/update
+72
View File
@@ -0,0 +1,72 @@
# AGENTS.md
> **Start every session by reading this file.**
> This file outlines the operational protocols, coding standards, and architectural decisions for this FastAPI project.
## 1. Agent Operational Protocols
### 🧠 Work Patterns (Plan-Act-Reflect)
* **Plan:** Before writing code, briefly outline your plan. Identify which files you will touch and what the side effects might be.
* **Act:** Execute the changes in small, atomic steps.
* **Reflect:** After coding, verify your work. Did you break existing tests? Did you add new tests?
### 🛡️ Git Discipline
* **NEVER commit to `main` or `master` directly.** Always create a feature branch: `feature/your-feature-name` or `fix/issue-description`.
* **Commit Messages:** Use the [Conventional Commits](https://www.conventionalcommits.org/) format.
* `feat: add user login endpoint`
* `fix: resolve database connection timeout`
* `refactor: split monolith dependency file`
* **Atomic Commits:** Keep commits small. One logical change = one commit.
### 📝 Changelog Maintenance
* **Update `CHANGELOG.md`** with every user-facing change.
* Format: `## [Unreleased] - YYYY-MM-DD` followed by `### Added`, `### Changed`, or `### Fixed`.
### 🚀 Release Flow
When changes are ready for deployment:
1. **Ask user if deploy cycle is desired **
2. **Update version** in `pyproject.toml`:
- Bug fixes: bump patch version (1.8.3 → 1.8.4)
- New features: bump minor version (1.8.4 → 1.9.0)
3. **Update CHANGELOG.md**:
- Move items from `[Unreleased]` to new version section
- Add release date: `## [1.8.4] - 2025-12-16`
4. **Commit and tag**:
```bash
git add -A
git commit -m "fix: description of changes"
git tag v1.8.4
git push origin main --tags
```
5. **CI/CD triggers automatically**:
- Gitea CI builds Docker image on new tag
- Watchtower pulls and deploys to production
- Verify deployment: `curl http://192.168.86.149:8000/health`
---
## 2. FastAPI Architecture & Best Practices
*Reference: [FastAPI Best Practices](https://github.com/zhanymkanov/fastapi-best-practices)*
### 📂 Project Structure (Directory-based, NOT File-type based)
Do **not** group files by type (e.g., one huge `routers` folder). Group by **domain/module** inside a `src/` directory.
**Correct Structure:**
```text
src/
├── auth/
│ ├── router.py # Endpoints
│ ├── schemas.py # Pydantic models
│ ├── service.py # Business logic (CRUD, etc.)
│ ├── dependencies.py# Module-specific dependencies
│ └── config.py # Module-specific settings
├── posts/
│ ├── router.py
│ └── ...
└── main.py # App entry point
+76
View File
@@ -4,6 +4,82 @@ All notable changes to The Scheduler will be documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/).
## [Unreleased]
## [1.4.0] - 2026-08-08
### Added
- **Portainer Backup Executor** (`portainer_backup_executor.py`) — archives Portainer's
own state through its `/api/backup` endpoint. Portainer's BoltDB lives in a Docker
volume that the daily config backup does not cover, so losing that volume would take
every stack definition with it. Uses the API rather than tarring the live volume, and
rejects a 200 whose body is not a readable archive.
## [1.3.0] - 2026-08-08
### Added
- **Postgres Retention Executor** (`postgres_retention_executor.py`) — deletes rows
past a retention window from a table on the shared Postgres server. Uses the
Scheduler's own credentials with only the database name overridden, so the target
database grants `scheduler_user` SELECT and DELETE on the table.
- **Docker Prune Executor** (`docker_prune_executor.py`) — scheduled reclaim of Docker
disk usage. Build cache and dangling images are pruned by default; unused images and
volumes are opt-in, since volume pruning also removes volumes belonging to stopped
containers.
### Fixed
- `POST /tasks` returned HTTP 500 after successfully creating the task. The response
model declared `created_at`/`updated_at` as strings while the database returns
timestamps, so every create looked like a failure and retrying hit a duplicate-key
error.
### Changed
- `TASK_REGISTRATION.md` now lists the executors that exist. It previously advertised
`shell`, `python` and `docker` executors that were never implemented.
## [1.2.0] - 2026-03-30
### Added
- **GCS Backup Executor** (`gcs_backup_executor.py`) — offsite backup to Google Cloud
Storage.
## [1.1.3] - 2026-01-08
### Changed
- Test release to validate CI/CD auto-deploy workflow
## [1.1.2] - 2026-01-03
### Fixed
- CI: Use curl for release creation (release-action requires Go)
## [1.1.1] - 2026-01-03
### Changed
- CI: Auto-create Gitea release on version tag push (v*) instead of manual release trigger
## [1.1.0] - 2025-12-14
### Added
- **Gitea Release Cleanup Executor** (`gitea_release_cleanup_executor.py`)
- Automatically cleans up old releases across all Gitea repositories
- Configurable retention count (default: 5 releases per repo)
- Repository exclusion list support
- Dry-run mode for safe testing
- Designed to run before Watchtower to prevent image tag accumulation
- **GITEA_TOKEN setting** in config for API token authentication (separate from password)
## [1.0.4] - 2025-12-14
### Fixed
- CI/CD: Correct Watchtower port (8080)
## [1.0.3] - 2025-12-14
### Added
- CI/CD: Trigger Watchtower update after successful Docker build
## [1.0.2] - 2025-12-14
### Fixed
+77 -5
View File
@@ -105,11 +105,83 @@ Calls HTTP endpoints. Supports environment variable substitution in headers/body
### Other Executors
- `shell`: Execute shell commands
- `python`: Execute Python scripts
- `docker`: Docker operations
- `backup`: Backup operations
- `doc_sync`: Documentation sync
The `executor` field is the module name under `src/executors/`. These are the
modules that actually exist:
- `config_backup_executor`: tar.gz backup of mounted directories, with retention
- `gcs_backup_executor`: offsite backup to Google Cloud Storage
- `doc_sync_executor`: mirror upstream docs into Gitea
- `gitea_release_cleanup_executor`: drop old Gitea releases, keeping the newest N
- `postgres_retention_executor`: delete rows past a retention window (see below)
- `docker_prune_executor`: reclaim Docker disk usage (see below)
- `portainer_backup_executor`: archive Portainer's own state via its backup API (see below)
- `example_executor`: demo/test
There is **no `shell` or `python` executor**. Earlier revisions of this document
listed them and they were never implemented; work needing a shell belongs either
in a purpose-built executor or on a host systemd timer.
#### `postgres_retention_executor`
Connects with the Scheduler's own Postgres credentials, overriding only the
database name, so the target database must grant `scheduler_user` SELECT and
DELETE on the table. Table and column names are validated against a strict
identifier pattern because they cannot be bound as query parameters.
```json
{
"database": "sysmon",
"table": "check_history",
"timestamp_column": "ts",
"retention_days": 30,
"dry_run": false
}
```
#### `portainer_backup_executor`
Portainer keeps every stack definition, endpoint, user and access-control rule in
a BoltDB inside the `portainer_data` Docker volume, which lives under
`/var/lib/docker/volumes/` and is **not** covered by the daily config backup.
This calls Portainer's `/api/backup` rather than tarring the volume: BoltDB is a
single memory-mapped file, so copying it live can capture a torn page.
The archive contains TLS certificates and private keys and is written `0600`. A
200 response whose body is not a readable archive is treated as a failure — an
archive that will not open is worse than a missing one, because it looks like a
backup until the day it is needed.
```json
{
"url": "${PORTAINER_URL}",
"api_key": "${PORTAINER_API_KEY}",
"output_dir": "/backups/portainer",
"retention_days": 30
}
```
Portainer runs host-networked, so a container name does not resolve; use the
host address. Requires `/mnt/media/backups/portainer` mounted into the container.
#### `docker_prune_executor`
Uses the docker socket already mounted into the container. Only the two stages
that discard regenerable data are on by default.
```json
{
"build_cache": true,
"dangling_images": true,
"unused_images": false,
"volumes": false,
"build_cache_until_hours": 168,
"dry_run": false
}
```
**`volumes` removes volumes belonging to merely-stopped containers, not just
orphaned ones.** Leave it off unless you have checked what is currently
unattached; on this host it is a plausible way to lose a database.
## Complete Task Schema
+3 -1
View File
@@ -1,6 +1,6 @@
[project]
name = "the-scheduler"
version = "1.0.2"
version = "1.4.0"
description = "System-wide maintenance orchestration - backups, doc mirroring, cleanup, task automation"
readme = "README.md"
requires-python = ">=3.12"
@@ -26,6 +26,8 @@ dependencies = [
"gitpython~=3.1.43",
# Security/Auth
"python-jose[cryptography]~=3.3.0",
# Cloud storage
"google-cloud-storage~=2.18.0",
]
[project.optional-dependencies]
+3
View File
@@ -26,6 +26,9 @@ gitpython~=3.1.43 # Latest stable
# Security/Auth
python-jose[cryptography]~=3.3.0 # JWT handling
# Cloud storage
google-cloud-storage~=2.18.0 # GCS offsite backups
# Testing
pytest~=8.3.4 # Test framework
pytest-asyncio~=0.25.2 # Async test support
+4
View File
@@ -57,6 +57,7 @@ class Settings(BaseSettings):
gitea_url: str = Field(default="http://gitea:3000", alias="GITEA_URL")
gitea_user: str = Field(default="library", alias="GITEA_USER")
gitea_password: str = Field(default="", alias="GITEA_PASSWORD")
gitea_token: str = Field(default="", alias="GITEA_TOKEN") # API token with write:repository scope
gitea_ssh_host: str = Field(default="gitea", alias="GITEA_SSH_HOST")
gitea_ssh_port: int = Field(default=22, alias="GITEA_SSH_PORT")
@@ -65,6 +66,9 @@ class Settings(BaseSettings):
backup_retention_weekly: int = Field(default=4, alias="BACKUP_RETENTION_WEEKLY")
backup_retention_monthly: int = Field(default=12, alias="BACKUP_RETENTION_MONTHLY")
# Google Cloud Storage
gcs_credentials_file: str = Field(default="", alias="GCS_CREDENTIALS_FILE")
# Documentation Mirroring
docs_mirror_path: str = Field(default="/docs-mirror", alias="DOCS_MIRROR_PATH")
docs_check_interval: int = Field(default=21600, alias="DOCS_CHECK_INTERVAL") # 6 hours
+109
View File
@@ -0,0 +1,109 @@
"""
Docker Prune Executor
Scheduled, non-interactive reclaim of Docker disk usage. The host equivalent is
system-admin-toj's scripts/disk/prune-docker.sh, which prompts per stage; a cron
task cannot prompt, so the destructive stages are opt-in instead.
Runs the docker CLI against the socket already mounted into this container.
Config schema:
{
"build_cache": true, # safe: cache is rebuilt on demand
"dangling_images": true, # safe: untagged layers nothing references
"unused_images": false, # re-pull on next deploy; costs bandwidth
"volumes": false, # DESTRUCTIVE - see below
"build_cache_until_hours": 168,
"dry_run": false
}
`volumes` is off by default and should stay off unless you have checked what is
actually unattached. `docker volume prune` removes every volume not bound to a
*running* container, which includes the data volume of anything merely stopped.
On this host that is a plausible way to lose a database.
Defaults are the two stages that only ever discard regenerable data.
"""
import asyncio
import logging
from typing import Any, Dict, List, Tuple
from src.config import Settings
logger = logging.getLogger(__name__)
COMMAND_TIMEOUT = 900
async def _run(args: List[str]) -> Tuple[int, str, str]:
proc = await asyncio.create_subprocess_exec(
*args,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
)
try:
stdout, stderr = await asyncio.wait_for(proc.communicate(), timeout=COMMAND_TIMEOUT)
except asyncio.TimeoutError:
proc.kill()
await proc.wait()
raise Exception(f"timed out after {COMMAND_TIMEOUT}s: {' '.join(args)}")
return proc.returncode, stdout.decode().strip(), stderr.decode().strip()
def _reclaimed(output: str) -> str:
"""Pull the 'Total reclaimed space: X' line out of docker's prune output."""
for line in output.splitlines():
if "reclaimed space" in line.lower():
return line.split(":", 1)[1].strip()
return "0B"
async def execute(config: Dict[str, Any], settings: Settings) -> str:
dry_run = bool(config.get("dry_run", False))
until_hours = int(config.get("build_cache_until_hours", 168))
stages: List[Tuple[str, List[str]]] = []
if config.get("build_cache", True):
stages.append(
("build cache", ["docker", "builder", "prune", "-f", "--filter", f"until={until_hours}h"])
)
if config.get("dangling_images", True):
stages.append(("dangling images", ["docker", "image", "prune", "-f"]))
if config.get("unused_images", False):
stages.append(("unused images", ["docker", "image", "prune", "-a", "-f"]))
if config.get("volumes", False):
logger.warning(
"volume pruning is enabled; this removes volumes belonging to stopped "
"containers, not just orphaned ones"
)
stages.append(("volumes", ["docker", "volume", "prune", "-f"]))
if not stages:
return "no prune stages enabled; nothing to do"
rc, out, err = await _run(["docker", "system", "df"])
if rc != 0:
raise Exception(f"docker unavailable: {err or out}")
before = out
if dry_run:
planned = ", ".join(name for name, _ in stages)
logger.info("dry run; would prune: %s", planned)
return f"dry run - would prune: {planned}\n{before}"
results = []
for name, args in stages:
rc, out, err = await _run(args)
if rc != 0:
# Report rather than abort: a later stage may still reclaim space, and
# a partial reclaim is more useful than none.
logger.error("prune stage %r failed: %s", name, err or out)
results.append(f"{name}: FAILED ({(err or out).splitlines()[0] if (err or out) else 'unknown'})")
continue
results.append(f"{name}: {_reclaimed(out)}")
logger.info("pruned %s -> %s", name, _reclaimed(out))
summary = "; ".join(results)
if any("FAILED" in r for r in results):
raise Exception(f"one or more prune stages failed: {summary}")
return f"reclaimed - {summary}"
+199
View File
@@ -0,0 +1,199 @@
"""
Google Cloud Storage Backup Executor
Backs up data to GCS buckets. Supports multiple modes:
- git_bundle: Creates a git bundle from a bare repo and uploads it
"""
import asyncio
import logging
import time
from datetime import datetime
from pathlib import Path
from typing import List, Optional
from google.cloud import storage
from src.config import Settings
logger = logging.getLogger(__name__)
async def execute(config: dict, settings: Settings) -> str:
"""
Execute GCS backup task.
Config schema:
{
"mode": "git_bundle",
"bucket": "bucket-name",
"prefix": "gitea/settled-reach",
"credentials_path": "/secrets/gcs-sa-key.json",
"repo_path": "/data/docker-data/gitea/data/git/repositories/user/repo.git",
"retention_count": 7,
"dry_run": false
}
Args:
config: Backup configuration
settings: Global scheduler settings
Returns:
Summary of backup operation
Raises:
Exception: On backup failure
"""
mode = config.get('mode')
if not mode:
raise ValueError("Missing required config field: mode")
if mode == 'git_bundle':
return await _mode_git_bundle(config, settings)
else:
raise ValueError(f"Unknown backup mode: {mode}")
async def _mode_git_bundle(config: dict, settings: Settings) -> str:
"""Create a git bundle from a bare repo and upload to GCS."""
bucket_name = config.get('bucket')
prefix = config.get('prefix', '').strip('/')
credentials_path = config.get('credentials_path', settings.gcs_credentials_file)
repo_path = Path(config.get('repo_path', ''))
retention_count = config.get('retention_count', 7)
dry_run = config.get('dry_run', False)
# Validate
if not bucket_name:
raise ValueError("Missing required config field: bucket")
if not repo_path or not str(repo_path).strip():
raise ValueError("Missing required config field: repo_path")
if not repo_path.exists():
raise ValueError(f"Repository path does not exist: {repo_path}")
if not (repo_path / 'HEAD').exists():
raise ValueError(f"Not a valid git repository: {repo_path}")
repo_name = repo_path.name.removesuffix('.git')
timestamp = datetime.now().strftime('%Y%m%d-%H%M%S')
bundle_filename = f"{repo_name}-{timestamp}.bundle"
bundle_path = Path(f"/tmp/{bundle_filename}")
logger.info(f"Starting GCS backup: mode=git_bundle, repo={repo_name}, dry_run={dry_run}")
try:
# Step 1: Create git bundle
logger.info(f"Creating git bundle from {repo_path}")
start = time.monotonic()
await _run_command([
'git', 'bundle', 'create',
str(bundle_path),
'--all'
], cwd=repo_path)
bundle_duration = time.monotonic() - start
if not bundle_path.exists():
raise Exception("Git bundle was not created")
bundle_size = bundle_path.stat().st_size
bundle_size_mb = bundle_size / (1024 * 1024)
logger.info(f"Bundle created: {bundle_filename} ({bundle_size_mb:.1f} MB) in {bundle_duration:.1f}s")
# Step 2: Verify bundle
await _run_command([
'git', 'bundle', 'verify',
str(bundle_path)
], cwd=repo_path)
logger.info("Bundle verified OK")
if dry_run:
return (
f"[DRY RUN] Would upload {bundle_filename} ({bundle_size_mb:.1f} MB) "
f"to gs://{bucket_name}/{prefix}/{bundle_filename}"
)
# Step 3: Upload to GCS
gcs_path = f"{prefix}/{bundle_filename}" if prefix else bundle_filename
logger.info(f"Uploading to gs://{bucket_name}/{gcs_path}")
start = time.monotonic()
client = storage.Client.from_service_account_json(credentials_path)
bucket = client.bucket(bucket_name)
blob = bucket.blob(gcs_path)
blob.upload_from_filename(str(bundle_path), timeout=3600)
upload_duration = time.monotonic() - start
logger.info(f"Upload complete in {upload_duration:.1f}s")
# Step 4: Retention cleanup in GCS
deleted_count = await _cleanup_gcs_retention(
client, bucket_name, prefix, repo_name, retention_count
)
return (
f"Backup completed: {bundle_filename} ({bundle_size_mb:.1f} MB). "
f"Bundle: {bundle_duration:.1f}s, Upload: {upload_duration:.1f}s. "
f"GCS: gs://{bucket_name}/{gcs_path}. "
f"Retention: {deleted_count} old bundle(s) removed."
)
finally:
# Always clean up the local temp file
if bundle_path.exists():
bundle_path.unlink()
logger.debug(f"Cleaned up temp file: {bundle_path}")
async def _cleanup_gcs_retention(
client: storage.Client,
bucket_name: str,
prefix: str,
repo_name: str,
retention_count: int
) -> int:
"""Delete old bundles from GCS, keeping only the most recent retention_count."""
if retention_count <= 0:
return 0
bucket = client.bucket(bucket_name)
blob_prefix = f"{prefix}/{repo_name}-" if prefix else f"{repo_name}-"
blobs = list(bucket.list_blobs(prefix=blob_prefix))
bundle_blobs = [b for b in blobs if b.name.endswith('.bundle')]
if len(bundle_blobs) <= retention_count:
logger.info(f"Retention OK: {len(bundle_blobs)} bundles (limit: {retention_count})")
return 0
# Sort by name (timestamp in name ensures chronological order)
bundle_blobs.sort(key=lambda b: b.name)
to_delete = bundle_blobs[:-retention_count]
for blob in to_delete:
logger.info(f"Deleting old bundle: {blob.name}")
blob.delete()
logger.info(f"Retention cleanup: deleted {len(to_delete)} old bundle(s)")
return len(to_delete)
async def _run_command(
cmd: List[str],
cwd: Optional[Path] = None,
) -> str:
"""Run shell command asynchronously."""
logger.debug(f"Running: {' '.join(cmd)} (cwd: {cwd})")
proc = await asyncio.create_subprocess_exec(
*cmd,
cwd=cwd,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE
)
stdout, stderr = await proc.communicate()
if proc.returncode != 0:
error_msg = stderr.decode() if stderr else "Unknown error"
raise Exception(f"Command failed: {' '.join(cmd)}\n{error_msg}")
return stdout.decode()
@@ -0,0 +1,211 @@
"""
Gitea Release Cleanup Executor
Cleans up old releases across all Gitea repositories, keeping only the most recent
releases per repository. Designed to run before Watchtower to prevent accumulation
of old container image tags.
Config schema:
{
"keep_count": 5, # Number of releases to keep per repo (default: 5)
"exclude_repos": [], # Repository names to skip (default: [])
"dry_run": false # If true, only log what would be deleted (default: false)
}
Example config:
{
"keep_count": 5,
"exclude_repos": ["important-repo", "legacy-app"],
"dry_run": false
}
Required settings:
- GITEA_URL: Base URL of Gitea instance
- GITEA_TOKEN: API token with write:repository scope
"""
import logging
from typing import Any
import httpx
from src.config import Settings
logger = logging.getLogger(__name__)
# Default configuration values
DEFAULT_KEEP_COUNT = 5
DEFAULT_TIMEOUT = 60
API_BASE = "/api/v1"
async def execute(config: dict, settings: Settings) -> str:
"""
Clean up old releases across all accessible Gitea repositories.
Args:
config: Task configuration (see module docstring)
settings: Global scheduler settings
Returns:
Summary of cleanup actions taken
Raises:
ValueError: On configuration error
Exception: On API call failure
"""
# Validate settings
if not settings.gitea_url:
raise ValueError("GITEA_URL not configured")
if not settings.gitea_token:
raise ValueError("GITEA_TOKEN not configured - required for release deletion")
# Parse configuration
keep_count = config.get("keep_count", DEFAULT_KEEP_COUNT)
exclude_repos = config.get("exclude_repos", [])
dry_run = config.get("dry_run", False)
if keep_count < 1:
raise ValueError(f"keep_count must be at least 1, got {keep_count}")
base_url = settings.gitea_url.rstrip("/")
headers = {"Authorization": f"token {settings.gitea_token}"}
stats = {
"repos_scanned": 0,
"repos_with_releases": 0,
"releases_deleted": 0,
"releases_skipped": 0,
"errors": [],
}
mode = "DRY RUN" if dry_run else "LIVE"
logger.info(f"Starting Gitea release cleanup ({mode}): keeping {keep_count} releases per repo")
async with httpx.AsyncClient(timeout=DEFAULT_TIMEOUT) as client:
# Fetch all repositories
repos = await _fetch_user_repos(client, base_url, headers)
logger.info(f"Found {len(repos)} repositories")
for repo in repos:
owner = repo["owner"]["login"]
name = repo["name"]
full_name = f"{owner}/{name}"
# Check exclusion list
if name in exclude_repos or full_name in exclude_repos:
logger.debug(f"Skipping excluded repo: {full_name}")
continue
stats["repos_scanned"] += 1
try:
# Fetch releases for this repo
releases = await _fetch_releases(client, base_url, headers, owner, name)
if not releases:
continue
stats["repos_with_releases"] += 1
# Determine which releases to delete (beyond keep_count)
to_delete = releases[keep_count:]
if not to_delete:
logger.debug(f"{full_name}: {len(releases)} releases, nothing to delete")
continue
logger.info(f"{full_name}: {len(releases)} releases, deleting {len(to_delete)}")
# Delete old releases
for release in to_delete:
release_id = release["id"]
tag_name = release["tag_name"]
if dry_run:
logger.info(f" [DRY RUN] Would delete: {tag_name} (id={release_id})")
stats["releases_skipped"] += 1
else:
try:
await _delete_release(client, base_url, headers, owner, name, release_id)
logger.info(f" Deleted: {tag_name}")
stats["releases_deleted"] += 1
except Exception as e:
error_msg = f"{full_name}/{tag_name}: {e}"
logger.warning(f" Failed to delete {tag_name}: {e}")
stats["errors"].append(error_msg)
except Exception as e:
error_msg = f"{full_name}: {e}"
logger.error(f"Error processing {full_name}: {e}")
stats["errors"].append(error_msg)
# Build summary
summary = _build_summary(stats, dry_run)
logger.info(f"Cleanup complete: {summary}")
return summary
async def _fetch_user_repos(
client: httpx.AsyncClient,
base_url: str,
headers: dict,
) -> list[dict[str, Any]]:
"""Fetch all repositories accessible to the authenticated user."""
url = f"{base_url}{API_BASE}/user/repos"
params = {"limit": 100}
response = await client.get(url, headers=headers, params=params)
response.raise_for_status()
return response.json()
async def _fetch_releases(
client: httpx.AsyncClient,
base_url: str,
headers: dict,
owner: str,
repo: str,
) -> list[dict[str, Any]]:
"""Fetch releases for a repository, sorted newest first (Gitea default)."""
url = f"{base_url}{API_BASE}/repos/{owner}/{repo}/releases"
params = {"limit": 100}
response = await client.get(url, headers=headers, params=params)
response.raise_for_status()
return response.json()
async def _delete_release(
client: httpx.AsyncClient,
base_url: str,
headers: dict,
owner: str,
repo: str,
release_id: int,
) -> None:
"""Delete a specific release."""
url = f"{base_url}{API_BASE}/repos/{owner}/{repo}/releases/{release_id}"
response = await client.delete(url, headers=headers)
response.raise_for_status()
def _build_summary(stats: dict, dry_run: bool) -> str:
"""Build a human-readable summary of the cleanup operation."""
parts = [
f"Scanned {stats['repos_scanned']} repos",
f"{stats['repos_with_releases']} with releases",
]
if dry_run:
parts.append(f"{stats['releases_skipped']} releases would be deleted")
else:
parts.append(f"{stats['releases_deleted']} releases deleted")
if stats["errors"]:
parts.append(f"{len(stats['errors'])} errors")
return ", ".join(parts)
+140
View File
@@ -0,0 +1,140 @@
"""
Portainer Backup Executor
Archives Portainer's own state through its `/api/backup` endpoint.
Why it needs backing up separately: Portainer keeps every stack definition,
endpoint, user and access-control rule in a BoltDB inside the Docker volume
`portainer_data`, which lives under /var/lib/docker/volumes/. The daily config
backup covers ~/docker-data and code-server-config only, so that volume is not
in it. Losing it takes all 24 stack definitions with it.
Why the API rather than tarring the volume: BoltDB is a single memory-mapped
file, so copying it while Portainer is writing can capture a torn page. The API
serialises a consistent snapshot.
The archive contains TLS certificates and private keys, so it is written 0600.
Config schema:
{
"url": "http://172.17.0.1:8001", # Portainer is host-networked, so a
# container name does not resolve;
# use the bridge gateway
"api_key": "${PORTAINER_API_KEY}", # ${VAR} reads the container env
"output_dir": "/backups/portainer",
"retention_days": 30,
"password": "" # optional; encrypts the archive
}
"""
import logging
import os
import re
import tarfile
import time
from datetime import datetime, timedelta, timezone
from pathlib import Path
import httpx
from src.config import Settings
logger = logging.getLogger(__name__)
BACKUP_TIMEOUT = 300
FILENAME_RE = re.compile(r"^portainer-\d{8}T\d{6}Z\.tar\.gz$")
def _substitute_env(value: str) -> str:
"""Expand ${VAR} against the container environment, as rest_api does."""
if not isinstance(value, str):
return value
for var in re.findall(r"\$\{([A-Z_][A-Z0-9_]*)\}", value):
resolved = os.getenv(var, "")
if not resolved:
logger.warning("environment variable not found: %s", var)
value = value.replace(f"${{{var}}}", resolved)
return value
def _prune(output_dir: Path, retention_days: int) -> int:
"""Delete archives older than the retention window. Returns how many went."""
cutoff = datetime.now(timezone.utc) - timedelta(days=retention_days)
removed = 0
for path in output_dir.glob("portainer-*.tar.gz"):
# Match the exact name this executor writes; never delete a stray file
# someone else put here.
if not FILENAME_RE.match(path.name):
continue
if datetime.fromtimestamp(path.stat().st_mtime, timezone.utc) < cutoff:
path.unlink()
removed += 1
logger.info("pruned old portainer backup: %s", path.name)
return removed
async def execute(config: dict, settings: Settings) -> str:
url = _substitute_env(config.get("url", "")).rstrip("/")
api_key = _substitute_env(config.get("api_key", ""))
output_dir = Path(config.get("output_dir", "/backups/portainer"))
retention_days = config.get("retention_days", 30)
password = _substitute_env(config.get("password", "") or "")
if not url:
raise ValueError("Missing required config: 'url'")
if not api_key:
raise ValueError("Missing or unresolved config: 'api_key'")
if not isinstance(retention_days, int) or isinstance(retention_days, bool) or retention_days < 1:
raise ValueError(f"retention_days must be a positive integer, got {retention_days!r}")
output_dir.mkdir(parents=True, exist_ok=True)
stamp = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ")
final = output_dir / f"portainer-{stamp}.tar.gz"
partial = final.with_suffix(".partial")
started = time.monotonic()
try:
async with httpx.AsyncClient(timeout=BACKUP_TIMEOUT) as client:
response = await client.post(
f"{url}/api/backup",
headers={"X-API-Key": api_key, "Content-Type": "application/json"},
json={"password": password} if password else {},
)
if response.status_code != 200:
raise Exception(
f"Portainer returned HTTP {response.status_code}: {response.text[:200]}"
)
partial.write_bytes(response.content)
# A 200 with a truncated body is still a failed backup. An archive that
# cannot be opened is worse than a missing one, because it looks like a
# backup until the day it is needed.
if not password:
try:
with tarfile.open(partial, "r:gz") as archive:
entries = len(archive.getnames())
except Exception as exc: # noqa: BLE001
# Deliberately broad. A truncated archive raises EOFError, which
# is neither TarError nor OSError, and any failure to open it
# means the same thing regardless of type: this is not a backup.
raise Exception(f"response is not a readable archive: {exc}") from exc
else:
entries = -1 # encrypted; contents cannot be verified here
partial.replace(final)
final.chmod(0o600) # contains TLS certs and private keys
finally:
if partial.exists():
partial.unlink()
removed = _prune(output_dir, retention_days)
kept = len([p for p in output_dir.glob("portainer-*.tar.gz") if FILENAME_RE.match(p.name)])
size_mb = final.stat().st_size / 1_048_576
elapsed = time.monotonic() - started
summary = (
f"backed up Portainer to {final.name} "
f"({size_mb:.2f} MB{'' if entries < 0 else f', {entries} entries'}, {elapsed:.1f}s); "
f"kept {kept}, pruned {removed} older than {retention_days}d"
)
logger.info(summary)
return summary
@@ -0,0 +1,109 @@
"""
Postgres Retention Executor
Deletes rows older than a retention window from a table on the shared Postgres
server. Written for sysmon's `check_history`, which grows with every monitoring
check and had no retention at all, but the executor is table-agnostic.
Connects with the Scheduler's own Postgres credentials and only overrides the
database name. That keeps a second set of credentials out of the stack; the
target database grants `scheduler_user` exactly SELECT and DELETE on the table,
so a bug here can drop old rows but cannot corrupt or forge history.
Config schema:
{
"database": "sysmon", # defaults to the Scheduler's own database
"table": "check_history", # required
"timestamp_column": "ts", # required
"retention_days": 30, # required, must be >= 1
"dry_run": false # count what would go, delete nothing
}
Table and column names cannot be passed as query parameters, so both are
validated against a strict identifier pattern before being interpolated.
Autovacuum reclaims the space afterwards; this deliberately does not VACUUM,
which would need table ownership the Scheduler intentionally does not have.
"""
import asyncio
import logging
import re
from typing import Any, Dict
import psycopg2
from src.config import Settings
logger = logging.getLogger(__name__)
# Deliberately strict: unquoted lowercase identifiers only. Anything needing
# quoting is out of scope and would be a hole in the interpolation below.
IDENTIFIER_RE = re.compile(r"^[a-z_][a-z0-9_]*$")
MAX_RETENTION_DAYS = 3650
def _validate_identifier(value: str, label: str) -> str:
if not isinstance(value, str) or not IDENTIFIER_RE.match(value):
raise ValueError(
f"invalid {label}: {value!r} (expected an unquoted lowercase identifier)"
)
return value
def _prune(config: Dict[str, Any], settings: Settings) -> str:
table = _validate_identifier(config.get("table", ""), "table")
column = _validate_identifier(config.get("timestamp_column", ""), "timestamp_column")
database = config.get("database") or settings.postgres_db
_validate_identifier(database, "database")
retention_days = config.get("retention_days")
if not isinstance(retention_days, int) or isinstance(retention_days, bool):
raise ValueError(f"retention_days must be an integer, got {retention_days!r}")
# A zero or negative window would delete everything, including the row the
# check just wrote. Refuse rather than quietly wipe the table.
if retention_days < 1 or retention_days > MAX_RETENTION_DAYS:
raise ValueError(
f"retention_days must be between 1 and {MAX_RETENTION_DAYS}, got {retention_days}"
)
dry_run = bool(config.get("dry_run", False))
cutoff_sql = f"{column} < now() - make_interval(days => %s)"
conn = psycopg2.connect(
host=settings.postgres_host,
port=settings.postgres_port,
database=database,
user=settings.postgres_user,
password=settings.postgres_password,
connect_timeout=10,
)
try:
with conn:
with conn.cursor() as cur:
cur.execute(f"SELECT count(*) FROM {table} WHERE {cutoff_sql}", (retention_days,))
stale = cur.fetchone()[0]
if dry_run:
logger.info("dry run: %s rows in %s.%s exceed %sd", stale, database, table, retention_days)
return f"dry run: {stale} rows older than {retention_days}d in {database}.{table}"
if stale == 0:
return f"nothing to prune in {database}.{table} (retention {retention_days}d)"
cur.execute(f"DELETE FROM {table} WHERE {cutoff_sql}", (retention_days,))
deleted = cur.rowcount
cur.execute(f"SELECT count(*) FROM {table}")
remaining = cur.fetchone()[0]
finally:
conn.close()
logger.info("pruned %s rows from %s.%s, %s remain", deleted, database, table, remaining)
return f"pruned {deleted} rows older than {retention_days}d from {database}.{table}, {remaining} remain"
async def execute(config: dict, settings: Settings) -> str:
"""Delete rows past the retention window. Returns a one-line summary."""
# psycopg2 is synchronous; keep it off the scheduler's event loop.
return await asyncio.to_thread(_prune, config, settings)
+8 -2
View File
@@ -3,6 +3,7 @@ Pydantic models for The Scheduler API.
"""
from pydantic import BaseModel, Field
from typing import Optional, Dict, Any
from datetime import datetime
from enum import Enum
@@ -203,8 +204,13 @@ class TaskUpdate(BaseModel):
class TaskResponse(TaskCreate):
"""Response model for task operations."""
created_at: str
updated_at: Optional[str] = None
# These are `timestamp` columns, so psycopg2 hands back datetime objects.
# Declaring them as `str` made Pydantic reject every create response, which
# 500'd the endpoint *after* the row had already been inserted and committed.
# FastAPI serialises datetime to an ISO 8601 string, so the JSON on the wire
# is unchanged — and now matches what GET /tasks/{name} already returned.
created_at: datetime
updated_at: Optional[datetime] = None
class Config:
from_attributes = True
+135
View File
@@ -0,0 +1,135 @@
"""
Tests for the docker prune executor.
The important property is which stages run. `volumes` removes volumes belonging
to merely-stopped containers, so it must never be enabled by accident, and the
safe stages must stay on by default.
"""
from unittest.mock import AsyncMock, patch
import pytest
from src.config import Settings
from src.executors import docker_prune_executor as prune
def _runner(reclaimed="Total reclaimed space: 1.5GB", rc=0):
"""Fake _run returning docker-shaped output for every invocation."""
async def run(args):
if args[:3] == ["docker", "system", "df"]:
return 0, "TYPE TOTAL ACTIVE SIZE RECLAIMABLE", ""
return rc, reclaimed, "" if rc == 0 else "boom"
return run
@pytest.mark.executor
@pytest.mark.unit
class TestStageSelection:
@pytest.mark.asyncio
async def test_defaults_run_only_the_safe_stages(self, test_settings: Settings):
calls = []
async def run(args):
calls.append(args)
if args[:3] == ["docker", "system", "df"]:
return 0, "df output", ""
return 0, "Total reclaimed space: 0B", ""
with patch.object(prune, "_run", run):
await prune.execute({}, test_settings)
joined = [" ".join(c) for c in calls]
assert any("builder prune" in c for c in joined)
assert any("image prune -f" in c for c in joined)
# The destructive ones must not appear without being asked for.
assert not any("volume prune" in c for c in joined)
assert not any("image prune -a" in c for c in joined)
@pytest.mark.asyncio
async def test_volumes_only_when_explicitly_enabled(self, test_settings: Settings):
calls = []
async def run(args):
calls.append(args)
if args[:3] == ["docker", "system", "df"]:
return 0, "df output", ""
return 0, "Total reclaimed space: 2GB", ""
with patch.object(prune, "_run", run):
await prune.execute({"volumes": True}, test_settings)
assert any("volume prune" in " ".join(c) for c in calls)
@pytest.mark.asyncio
async def test_all_stages_disabled_is_a_no_op(self, test_settings: Settings):
with patch.object(prune, "_run", AsyncMock()) as run:
result = await prune.execute(
{"build_cache": False, "dangling_images": False}, test_settings
)
assert "nothing to do" in result
run.assert_not_called()
@pytest.mark.asyncio
async def test_dry_run_executes_no_prune(self, test_settings: Settings):
calls = []
async def run(args):
calls.append(args)
return 0, "df output", ""
with patch.object(prune, "_run", run):
result = await prune.execute({"dry_run": True}, test_settings)
assert "dry run" in result
assert all("prune" not in " ".join(c) for c in calls)
@pytest.mark.executor
@pytest.mark.unit
class TestFailureHandling:
@pytest.mark.asyncio
async def test_docker_unavailable_raises(self, test_settings: Settings):
async def run(args):
return 1, "", "Cannot connect to the Docker daemon"
with patch.object(prune, "_run", run):
with pytest.raises(Exception, match="docker unavailable"):
await prune.execute({}, test_settings)
@pytest.mark.asyncio
async def test_failed_stage_surfaces_but_others_still_run(self, test_settings: Settings):
attempted = []
async def run(args):
if args[:3] == ["docker", "system", "df"]:
return 0, "df output", ""
attempted.append(" ".join(args))
if "builder" in args:
return 1, "", "builder exploded"
return 0, "Total reclaimed space: 3MB", ""
with patch.object(prune, "_run", run):
with pytest.raises(Exception, match="one or more prune stages failed"):
await prune.execute({}, test_settings)
# The image stage must still have been attempted after builder failed.
assert any("image prune" in a for a in attempted)
@pytest.mark.executor
@pytest.mark.unit
class TestOutputParsing:
@pytest.mark.parametrize(
"output,expected",
[
("Total reclaimed space: 1.5GB", "1.5GB"),
("deleted: sha256:abc\nTotal reclaimed space: 0B", "0B"),
("no such line", "0B"),
("", "0B"),
],
)
def test_reclaimed_parsing(self, output, expected):
assert prune._reclaimed(output) == expected
+175
View File
@@ -0,0 +1,175 @@
"""
Tests for the Portainer backup executor.
The point of this executor is producing an archive that will still open on the
day it is needed, so most of these cover the failure paths: a truncated body
behind a 200, a partial file left on disk, and retention deleting the wrong
thing.
"""
import gzip
import io
import tarfile
from datetime import datetime, timedelta, timezone
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
from src.config import Settings
from src.executors import portainer_backup_executor as pbe
def _tar_gz_bytes(names=("compose/1/docker-compose.yml", "certs/cert.pem")) -> bytes:
buf = io.BytesIO()
with tarfile.open(fileobj=buf, mode="w:gz") as tar:
for n in names:
data = b"x"
info = tarfile.TarInfo(name=n)
info.size = len(data)
tar.addfile(info, io.BytesIO(data))
return buf.getvalue()
def _mock_post(status=200, content=None):
response = MagicMock()
response.status_code = status
response.content = content if content is not None else _tar_gz_bytes()
response.text = "error body"
client = MagicMock()
client.post = AsyncMock(return_value=response)
ctx = MagicMock()
ctx.__aenter__ = AsyncMock(return_value=client)
ctx.__aexit__ = AsyncMock(return_value=False)
return ctx, client
@pytest.mark.executor
@pytest.mark.unit
class TestConfigValidation:
@pytest.mark.asyncio
async def test_missing_url_rejected(self, test_settings: Settings, tmp_path):
with pytest.raises(ValueError, match="url"):
await pbe.execute({"api_key": "k", "output_dir": str(tmp_path)}, test_settings)
@pytest.mark.asyncio
async def test_missing_api_key_rejected(self, test_settings: Settings, tmp_path):
with pytest.raises(ValueError, match="api_key"):
await pbe.execute({"url": "http://x", "output_dir": str(tmp_path)}, test_settings)
@pytest.mark.asyncio
async def test_unresolved_env_var_rejected(self, test_settings: Settings, tmp_path, monkeypatch):
"""${VAR} that expands to nothing must fail, not send an empty key."""
monkeypatch.delenv("NOPE_MISSING", raising=False)
with pytest.raises(ValueError, match="api_key"):
await pbe.execute(
{"url": "http://x", "api_key": "${NOPE_MISSING}", "output_dir": str(tmp_path)},
test_settings,
)
@pytest.mark.asyncio
@pytest.mark.parametrize("bad", [0, -1, "30", None, True])
async def test_bad_retention_rejected(self, bad, test_settings: Settings, tmp_path):
with pytest.raises(ValueError, match="retention_days"):
await pbe.execute(
{"url": "http://x", "api_key": "k", "output_dir": str(tmp_path),
"retention_days": bad},
test_settings,
)
@pytest.mark.executor
@pytest.mark.unit
class TestBackupBehaviour:
def _config(self, tmp_path, **over):
cfg = {"url": "http://portainer:9000", "api_key": "k",
"output_dir": str(tmp_path), "retention_days": 30}
cfg.update(over)
return cfg
@pytest.mark.asyncio
async def test_writes_verified_archive(self, test_settings: Settings, tmp_path):
ctx, _ = _mock_post()
with patch("httpx.AsyncClient", return_value=ctx):
result = await pbe.execute(self._config(tmp_path), test_settings)
files = list(tmp_path.glob("portainer-*.tar.gz"))
assert len(files) == 1
assert "2 entries" in result
with tarfile.open(files[0], "r:gz") as tar: # opens = usable backup
assert "certs/cert.pem" in tar.getnames()
@pytest.mark.asyncio
async def test_archive_is_not_world_readable(self, test_settings: Settings, tmp_path):
"""It contains TLS private keys."""
ctx, _ = _mock_post()
with patch("httpx.AsyncClient", return_value=ctx):
await pbe.execute(self._config(tmp_path), test_settings)
f = next(tmp_path.glob("portainer-*.tar.gz"))
assert oct(f.stat().st_mode)[-3:] == "600"
@pytest.mark.asyncio
async def test_http_error_raises_and_leaves_nothing(self, test_settings: Settings, tmp_path):
ctx, _ = _mock_post(status=401)
with patch("httpx.AsyncClient", return_value=ctx):
with pytest.raises(Exception, match="HTTP 401"):
await pbe.execute(self._config(tmp_path), test_settings)
assert list(tmp_path.iterdir()) == []
@pytest.mark.asyncio
async def test_truncated_body_behind_200_is_rejected(self, test_settings: Settings, tmp_path):
"""The dangerous case: a 200 whose body is not a usable archive."""
broken = _tar_gz_bytes()[:40]
ctx, _ = _mock_post(content=broken)
with patch("httpx.AsyncClient", return_value=ctx):
with pytest.raises(Exception, match="not a readable archive"):
await pbe.execute(self._config(tmp_path), test_settings)
# no .partial and no final file left behind
assert list(tmp_path.iterdir()) == []
@pytest.mark.asyncio
async def test_gzip_that_is_not_a_tar_is_rejected(self, test_settings: Settings, tmp_path):
ctx, _ = _mock_post(content=gzip.compress(b"not a tar"))
with patch("httpx.AsyncClient", return_value=ctx):
with pytest.raises(Exception, match="not a readable archive"):
await pbe.execute(self._config(tmp_path), test_settings)
assert list(tmp_path.iterdir()) == []
@pytest.mark.asyncio
async def test_api_key_resolved_from_env(self, test_settings: Settings, tmp_path, monkeypatch):
monkeypatch.setenv("PT_KEY", "secret-value")
ctx, client = _mock_post()
with patch("httpx.AsyncClient", return_value=ctx):
await pbe.execute(self._config(tmp_path, api_key="${PT_KEY}"), test_settings)
assert client.post.call_args.kwargs["headers"]["X-API-Key"] == "secret-value"
@pytest.mark.executor
@pytest.mark.unit
class TestRetention:
def _age(self, path, days):
import os
old = (datetime.now(timezone.utc) - timedelta(days=days)).timestamp()
os.utime(path, (old, old))
def test_prunes_only_past_the_window(self, tmp_path):
fresh = tmp_path / "portainer-20260808T120000Z.tar.gz"
stale = tmp_path / "portainer-20260101T120000Z.tar.gz"
for f in (fresh, stale):
f.write_bytes(b"x")
self._age(stale, 45)
assert pbe._prune(tmp_path, 30) == 1
assert fresh.exists() and not stale.exists()
def test_leaves_unrelated_files_alone(self, tmp_path):
"""Retention must not touch anything it did not write."""
other = tmp_path / "important-database-dump.tar.gz"
named_alike = tmp_path / "portainer-backup-manual.tar.gz"
for f in (other, named_alike):
f.write_bytes(b"x")
self._age(f, 400)
assert pbe._prune(tmp_path, 30) == 0
assert other.exists() and named_alike.exists()
+151
View File
@@ -0,0 +1,151 @@
"""
Tests for the postgres retention executor.
Focus is on the guards. The executor interpolates a table and column name
straight into SQL (they cannot be bound as parameters), and it issues DELETEs
against a live table, so the validation in front of both is what keeps a
malformed config from becoming data loss.
"""
from unittest.mock import MagicMock, patch
import pytest
from src.config import Settings
from src.executors import postgres_retention_executor as retention
@pytest.mark.executor
@pytest.mark.unit
class TestIdentifierValidation:
"""Table/column/database names are interpolated, so they must be rejected early."""
@pytest.mark.parametrize(
"bad",
[
"check_history; DROP TABLE users",
'check_history"',
"check history",
"Check_History", # uppercase would need quoting to resolve
"1_history",
"",
"--comment",
],
)
def test_rejects_unsafe_identifiers(self, bad):
with pytest.raises(ValueError):
retention._validate_identifier(bad, "table")
@pytest.mark.parametrize("good", ["check_history", "ts", "_private", "a1"])
def test_accepts_plain_identifiers(self, good):
assert retention._validate_identifier(good, "table") == good
@pytest.mark.executor
@pytest.mark.unit
class TestRetentionGuards:
"""A bad retention window must never reach the database."""
def _config(self, **overrides):
config = {
"database": "sysmon",
"table": "check_history",
"timestamp_column": "ts",
"retention_days": 30,
}
config.update(overrides)
return config
@pytest.mark.parametrize("days", [0, -1, -30, 3651])
def test_rejects_out_of_range_retention(self, days, test_settings: Settings):
# 0 or negative would delete every row including the one just written.
with patch("psycopg2.connect") as connect:
with pytest.raises(ValueError):
retention._prune(self._config(retention_days=days), test_settings)
connect.assert_not_called()
@pytest.mark.parametrize("days", ["30", None, 1.5, True])
def test_rejects_non_integer_retention(self, days, test_settings: Settings):
with patch("psycopg2.connect") as connect:
with pytest.raises(ValueError):
retention._prune(self._config(retention_days=days), test_settings)
connect.assert_not_called()
def test_rejects_injection_in_table_before_connecting(self, test_settings: Settings):
with patch("psycopg2.connect") as connect:
with pytest.raises(ValueError):
retention._prune(
self._config(table="check_history; DELETE FROM check_history --"),
test_settings,
)
connect.assert_not_called()
@pytest.mark.executor
@pytest.mark.unit
class TestRetentionBehaviour:
"""Behaviour against a mocked cursor."""
def _mock_conn(self, counts):
cursor = MagicMock()
cursor.fetchone.side_effect = [(c,) for c in counts]
cursor.rowcount = counts[0] if counts else 0
conn = MagicMock()
conn.cursor.return_value.__enter__.return_value = cursor
conn.__enter__.return_value = conn
return conn, cursor
def _config(self, **overrides):
config = {
"database": "sysmon",
"table": "check_history",
"timestamp_column": "ts",
"retention_days": 30,
}
config.update(overrides)
return config
def test_dry_run_does_not_delete(self, test_settings: Settings):
conn, cursor = self._mock_conn([7])
with patch("psycopg2.connect", return_value=conn):
result = retention._prune(self._config(dry_run=True), test_settings)
assert "dry run" in result
assert "7" in result
executed = " ".join(str(c) for c in cursor.execute.call_args_list)
assert "DELETE" not in executed.upper()
def test_no_stale_rows_skips_delete(self, test_settings: Settings):
conn, cursor = self._mock_conn([0])
with patch("psycopg2.connect", return_value=conn):
result = retention._prune(self._config(), test_settings)
assert "nothing to prune" in result
executed = " ".join(str(c) for c in cursor.execute.call_args_list)
assert "DELETE" not in executed.upper()
def test_deletes_and_reports(self, test_settings: Settings):
# count(stale) -> 5, then count(remaining) -> 42
conn, cursor = self._mock_conn([5, 42])
cursor.rowcount = 5
with patch("psycopg2.connect", return_value=conn):
result = retention._prune(self._config(), test_settings)
assert "pruned 5 rows" in result
assert "42 remain" in result
executed = " ".join(str(c) for c in cursor.execute.call_args_list)
assert "DELETE" in executed.upper()
def test_defaults_to_scheduler_database(self, test_settings: Settings):
conn, _ = self._mock_conn([0])
config = self._config()
del config["database"]
with patch("psycopg2.connect", return_value=conn) as connect:
retention._prune(config, test_settings)
assert connect.call_args.kwargs["database"] == test_settings.postgres_db
@pytest.mark.asyncio
async def test_execute_wraps_prune(self, test_settings: Settings):
conn, _ = self._mock_conn([0])
with patch("psycopg2.connect", return_value=conn):
result = await retention.execute(self._config(), test_settings)
assert "nothing to prune" in result