29 Commits
Author SHA1 Message Date
jpmschweitzerandClaude d91d0f2c63 release v1.5.1
Build and Push / release (push) Successful in 2s
Build and Push / build (push) Successful in 1m14s
Also collapses a duplicated changelog section. A second "## [Unreleased]"
heading has existed at the foot of the file since 2026-08-07, above the
roadmap list -- and because it matched first, my own v1.5.0 and T-69 commits
wrote their entries into both copies. The roadmap is now "## Planned", which is
what it always was, so the two cannot collide again.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-11 11:35:29 +02:00
jpmschweitzerandClaude b62024b21d fix(executors): run backups off the event loop
The scheduler serves its own REST API from the same loop that runs executors,
and the config backup spends ~21 minutes inside tarfile and zlib. Called inline
that starves the loop for the whole window: the service was unreachable
03:05-03:25 every night, and again at 07:39 today when the job was triggered by
hand to prove the T-69 report path. asyncio.to_thread is the fix -- zlib
releases the GIL while compressing, so the loop is scheduled normally.

The outage was invisible for as long as it existed. The hourly health check
fires at :35 and the outage runs 03:05-03:25, so no sample ever landed inside
it. A fixed-phase hourly probe cannot see a 20-minute event; that is aliasing,
not bad luck, and it would have stayed hidden indefinitely.

health_report.report_async joins it: psycopg2 is a blocking driver, so
reporting from the loop held it for the connect and insert -- up to
connect_timeout seconds precisely when the database is unreachable, which is
when a report matters most. Both backup executors use it now.

Caveat worth knowing: the engine wraps executors in asyncio.wait_for and a
thread cannot be cancelled, so on timeout the task is recorded failed while the
tar runs to completion. Still strictly better than blocking everything, and the
configured 3600s is well clear of the observed 1263s.

Seven tests in this file had been red since the initial commit -- they came
over with the portainer-core extraction, patched the Path class wholesale,
asserted "backed up" against a function returning "Backup completed: ...", and
one wrapped its call in except Exception: pass with its only assertion
commented out. There is no CI test gate here, so nothing reported it. Replaced
with tests that build real archives in tmp_path and assert on their contents.

The new loop test is the one that matters and it is mutation-checked: with
to_thread reverted it counts 0 heartbeat ticks, with it ~40.

Their structural demands (_create_tar_filter, a sync _cleanup_old_backups) were
adopted because threading wanted that shape anyway. Their exclude semantics
were not. The tests assert fnmatch behaviour and the deployed config is written
against substring matching -- it excludes logs as ".log", and paths as
unanchored fragments like "ollama/models/*" against members named
"docker-data/ollama/...". Under fnmatch neither matches, and the nightly
archive would silently gain many GB of model blobs instead of shrinking. That
is now pinned by tests naming the consequence, so the "improvement" fails loudly.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-11 11:35:29 +02:00
jpmschweitzerandClaude 18be804ba0 fix: report the database's refusal instead of guessing its cause
The health report's InsufficientPrivilege handler said 'the scheduler's
database user lacks INSERT on check_history'. That grant was in place. What was
actually missing was USAGE on check_history_id_seq — the sequence behind the
table's serial id — so an INSERT was refused for a reason the message did not
mention and actively contradicted.

Verified by attempting the insert directly as scheduler_user:
  ERROR: permission denied for sequence check_history_id_seq

A diagnostic that names a cause it did not observe is worse than a generic one:
it sends the reader to a fix that is already applied, and reads as evidence the
grant did not work. The message now prints psycopg2's own first line.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-11 10:04:41 +02:00
jpmschweitzerandClaude 1012f6f374 release v1.5.0
Build and Push / release (push) Successful in 3s
Build and Push / build (push) Successful in 1m14s
Backup executors report their own outcome to check_history. The grant that
makes it functional — INSERT for scheduler_user on the sysmon database — was
applied 2026-08-11, so this deploy is the step that closes the gap: the backup
domain has had no writer since the file-age poll was retired.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-11 09:17:27 +02:00
jpmschweitzer d9cbaee1fc feat: backup executors report their own outcome to check_history
The health record used to learn about backups by polling the mtime of the
newest archive, hourly, against a 48-hour threshold — for a job that runs once
a day. Forty-seven of every forty-eight runs could not produce a new answer.

Worse, file age cannot distinguish a failed backup from one that has not run
yet: when tonight's job dies, yesterday's archive is 24 hours old and still
reads healthy, and keeps reading healthy until hour 48. A failure stayed
invisible for two days to the check whose only job was noticing it. These
executors know at 03:05.

Reporting wraps the work rather than living inside it, so the failure path
cannot be forgotten — an exception is reported as critical and re-raised,
leaving the task's own status untouched. Reporting only success would reproduce
exactly the blind spot this replaces.

A reporting failure never fails the backup. Everything in health_report is
caught and logged, which means the absence of rows is the only symptom a broken
reporter produces — so monitor for rows, not for errors.

Not yet functional in production: scheduler_user holds DELETE and SELECT on
check_history (it prunes the table nightly) but not INSERT. Until that grant is
made the report logs a refusal and skips. The ticket predicted this as the step
that would fail silently, which is why it is named in the changelog and handled
as its own exception rather than folded into a generic catch.

Verified against baseline: tests/test_config_backup_executor.py and
tests/test_portainer_backup_executor.py report 7 failed / 17 passed both with
and without this change, so the pre-existing failures are untouched.

Workspace D-33, T-69.
2026-08-11 00:11:52 +02:00
jpmschweitzerandClaude 0c2199667f build(ci): move the pre-push gate into the Makefile
The hook carried ~50 lines of gitleaks logic and a comment explaining it was
self-contained because "this repo has no Makefile". It has one now, so the
reason is gone and the arrangement is backwards: a hook is a trigger, and
logic belongs where it can be read, run by hand, and changed under review.

.githooks/pre-push is now a byte-identical shim onto `make pre-push` in every
repo in the workspace. The scan itself moves to ci/secrets.sh unchanged, and
`make secrets` runs it on its own.

The call surface is identical everywhere; what it runs is not, and should not
be — each repo gates what it actually has. That is the point of standardising
the name rather than the contents: nobody has to read a repo to find out how
to check it.

secrets runs first, deliberately. It is the only failure here that cannot be
undone by fixing it afterwards — a failed lint costs another commit, a pushed
credential is cached and indexed whether or not it is later deleted.

Some of these gates fail today, on lint debt that predates them, and they are
left wired anyway. The board was measured once and written down in T-56
instead of being worked around here. Narrowing each gate to whatever already
passes would produce a gate that reports success for doing nothing, which is
the failure this workspace keeps rediscovering.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-09 18:57:22 +02:00
jpmschweitzerandClaude 054974c3fb ci(make): reserve exit 69 for "could not run" (D-26)
Environment guards now exit 69 rather than 1, so a caller can tell a suite
that could not start from one that ran and failed. The first toj test sweep
reported "3 repositories failed" and none of the three had executed a test —
two could not find go, one had no venv. That points the reader at the tests
when the fault is in the environment.

Only the environment guards change. A gitleaks finding, a failed test run and
a vulncheck hit still exit 1, because those did run and did fail.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-09 15:56:21 +02:00
jpmschweitzerandClaude b712caefb1 build: add the Makefile command surface (D-27)
Every repo gets one at the root: help, plus test and lint where those exist.
The point is that a target name means the same thing in every repo, so an
agent or a person can act without reading the repo first.

Paths resolve here rather than in callers (D-10). python3 on this host is 3.8
and cannot parse these sources, and a bare pytest or ruff resolves only in a
login shell — so both are named explicitly through the venv, and a missing
venv fails with the command to fix it rather than a bare no-such-file.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-09 15:10:33 +02:00
jpmschweitzerandClaude 5b4814a3cd chore(claude): pin PQL_VAULT per project so cwd stops choosing the vault
pql is now a bare word on PATH, which removed the long incantation that had
been forcing --vault into every call by habit. Convenience lowered the cost
of the wrong thing without lowering the cost of the right one: a three-word
pql ticket new targets whichever vault the cwd happens to sit in, and there
are nine of them with colliding id sequences.

PQL_VAULT in each project settings file makes the vault a property of the
session rather than of the working directory — the same lesson Rule 3 records
for git -C, applied to pql. Verified the env var overrides cwd discovery,
that an explicit --vault still beats the env var, and that the harness
hot-reloads it without a restart.

This does not make provenance visible: no output says which vault answered,
so a forgotten --vault still returns a well-formed answer about the wrong
dataset. That remains T-37.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-09 13:49:00 +02:00
jpmschweitzerandClaude 9eb9e4b32b chore(claude): deny toj in the sub-repos
toj is now on the global PATH as /usr/local/bin/toj, so its scope boundary
had to stop being "the absolute path is inconvenient to type" and start
being a rule. Its repo and settings verbs operate on the workspace root; run
from inside this repo they answer about the wrong tree.

Both spellings are denied, bare and absolute, because a deny with one
spelling left open is decorative.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-09 13:42:06 +02:00
jpmschweitzerandClaude c64412d9b3 ci: gate pushes on a gitleaks scan of the outgoing commits
No repo here scanned for committed credentials. The hook is self-contained
rather than delegating to a Makefile, because this repo has none and a hook
reaching into a sibling repo breaks the moment this one is cloned elsewhere.

Scans the outgoing range rather than full history: history carries settled
findings — test fixtures, vendored third-party code — and a gate that fails
on something unfixable gets bypassed within a week.

Setting core.hooksPath means pql init must replant its replication shims into
.githooks, which is why they are gitignored here alongside the tracked
pre-push. Same layout pql itself uses.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-09 12:48:54 +02:00
jpmschweitzerandClaude 4cbfe7a6f9 docs: qualify workspace decision ids cited from this repo
Decision ids are per-vault sequences, so they collide by construction
once there is more than one vault -- and every repo now has one. A bare
D-15 here will mean this repo's D-15 the moment this repo records one.
Cross-vault references are therefore qualified: workspace D-15.

Not hypothetical: pql holds D-1 through D-31 while the workspace holds
D-1 through D-21, so every workspace id currently collides with an
unrelated pql one. A bare id is not wrong the day it is written -- it
decays into wrong as the other vault grows, and nothing flags it.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-09 04:17:07 +02:00
jpmschweitzerandClaude e975f5b720 docs: replace AGENTS.md with a repo-specific CLAUDE.md
One agent doc per repo, and it is CLAUDE.md. Written fresh rather than
reformatted. The old file's feature-branch mandate and its `git add -A`
release snippet are both gone, and its health-check URL pointed at port
8000, which is tatlock -- this service is 8090.

The section worth reading is on executors, because neither of the usual
ways to establish what code is live works here. src/executors/*.py are
never statically imported: src/tasks/executor.py builds the module path
from a scheduled_tasks row and calls __import__ at execution time. So a
grep finds no importer, and a cold sys.modules snapshot shows none of
them loaded. The authoritative source is the database, and the doc
carries the query.

That distinction matters for gcs_backup_executor, which has zero rows
today. It is dormant, not dead: it becomes live the moment someone
inserts a row naming it, with no code change and no deploy.

Also records that the table is scheduled_tasks even though the API path
is /tasks, so the obvious query fails with UndefinedTable, and that the
empty src/config/ directory does not shadow src/config.py -- verified in
the container, a regular module wins over a namespace package.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-09 03:15:15 +02:00
jpmschweitzerandClaude 55fe0f9083 chore: adopt the workspace agent-config baseline
Commits a .claude/settings.json rather than leaving permissions to
per-developer local state, and initialises a pql vault for this repo's
tickets and internal decisions.

Every git deny rule appears in both the `git <verb>` and `git * <verb>`
forms. Only the second catches `git -C <path>`, and without it the whole
deny list is decorative -- it looks like a policy and stops nothing.

The allow list carries pql's absolute path alongside the bare name.
pql is installed to ~/.local/bin, which is on the login PATH but not the
one a non-interactive shell gets, so the bare-name rules match nothing on
their own and every call would prompt anyway.

.gitignore now covers .claude/settings.local.json, which is machine-local
and must never be shared. `pql init` contributed the .pql/* rules with an
exception for the changelog, which is the replication log of record and
has to be committed for tickets to travel with a clone.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-09 03:14:59 +02:00
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
32 changed files with 3058 additions and 249 deletions
+70
View File
@@ -0,0 +1,70 @@
{
"env": {
"PQL_VAULT": "/mnt/media/Projects/scheduler"
},
"permissions": {
"allow": [
"Bash(pql)",
"Bash(pql *)",
"Bash(/home/jpmschweitzer/.local/bin/pql:*)",
"Bash(git status:*)",
"Bash(git log:*)",
"Bash(git diff:*)",
"Bash(git branch:*)",
"Bash(.venv/bin/python -m pytest:*)",
"Bash(.venv/bin/pytest:*)",
"Bash(pytest:*)",
"Bash(docker logs scheduler:*)",
"Bash(curl -s http://localhost:8090/*)"
],
"deny": [
"Bash(/mnt/media/Projects/cladmin/ops/bin/toj)",
"Bash(/mnt/media/Projects/cladmin/ops/bin/toj:*)",
"Bash(chmod -R 777 *)",
"Bash(chmod 777 *)",
"Bash(dd if=*)",
"Bash(find * -delete*)",
"Bash(find * -exec*)",
"Bash(git * add --all*)",
"Bash(git * add -A*)",
"Bash(git * add .)",
"Bash(git * branch -D *)",
"Bash(git * checkout -- *)",
"Bash(git * clean -fd*)",
"Bash(git * clean -fdx*)",
"Bash(git * commit --no-verify*)",
"Bash(git * merge --no-ff*)",
"Bash(git * push --force*)",
"Bash(git * push -f*)",
"Bash(git * reset --hard*)",
"Bash(git * restore .*)",
"Bash(git add --all*)",
"Bash(git add -A*)",
"Bash(git add .)",
"Bash(git branch -D *)",
"Bash(git checkout -- *)",
"Bash(git clean -fd*)",
"Bash(git clean -fdx*)",
"Bash(git commit --no-verify*)",
"Bash(git merge --no-ff*)",
"Bash(git push --force*)",
"Bash(git push -f*)",
"Bash(git reset --hard*)",
"Bash(git restore .*)",
"Bash(mkfs*)",
"Bash(psql * -c DELETE FROM scheduled_tasks*)",
"Bash(psql * -c DROP*)",
"Bash(psql * DROP DATABASE*)",
"Bash(psql * TRUNCATE*)",
"Bash(redis-cli * FLUSHALL*)",
"Bash(redis-cli * FLUSHDB*)",
"Bash(rm -rf $HOME*)",
"Bash(rm -rf /*)",
"Bash(rm -rf ~*)",
"Bash(su *)",
"Bash(sudo *)",
"Bash(toj)",
"Bash(toj:*)"
]
}
}
+1
View File
@@ -0,0 +1 @@
.pql/changelog/*.sql merge=union
+18 -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,8 +36,8 @@ 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()
+13
View File
@@ -0,0 +1,13 @@
#!/usr/bin/env bash
# Trigger only. The checks live in the Makefile, where they can be read, run by
# hand (`make pre-push`), and changed under review.
#
# This file is identical in every repo in this workspace, deliberately: the call
# surface is the same everywhere even though what each gate runs is not, so
# nobody has to read a repo to find out how to check it (D-27).
#
# Enable per clone with: git config core.hooksPath .githooks
# Never bypass with --no-verify. Suppress a specific finding deliberately
# instead, with a reason — see `make pre-push`.
set -euo pipefail
exec make -C "$(git rev-parse --show-toplevel)" pre-push
+13
View File
@@ -87,3 +87,16 @@ cython_debug/
# Project-specific
logs/
task-data/
# Claude Code user-specific settings
.claude/settings.local.json
.pql/*
!.pql/changelog/
# pql shims planted by `pql init` into the dir core.hooksPath points at.
# Per-clone: each embeds the absolute path of the pql binary that planted it.
# Only .githooks/pre-push is shared.
.githooks/pre-commit
.githooks/post-merge
.githooks/post-checkout
.githooks/post-rewrite
+11
View File
@@ -0,0 +1,11 @@
-- Changelog format marker, written by pql. Comments only: this file
-- is never executed — Import descends into the per-table directories
-- and does not read the changelog root.
--
-- A changelog carrying no marker is format 1, the shape that existed
-- before formats were versioned. An older format is migrated forward
-- by `pql plan upgrade` (and automatically from the post-merge hook);
-- a newer one is refused rather than replayed under rules this binary
-- does not know. See D-28 and docs/versions.md.
-- pql:changelog_format: 2.0.0
-- pql:written_by: 2.2.0
+139
View File
@@ -0,0 +1,139 @@
-- Auto-generated by pql init. CREATE TABLE statements
-- for the planning schema; per-table dir keeps the changelog
-- self-describing per D-15. CREATE TABLE IF NOT EXISTS is
-- idempotent so running schema files from each directory in
-- replay order is harmless.
--
-- Importer parses the markers below to detect schema drift
-- between the producing pql version and the local one — a
-- bumped canonical_version means projection rules changed
-- and replay must refuse rather than silently corrupt state.
-- pql:created_by: 2.2.0
-- pql:canonical_version: 2
CREATE TABLE IF NOT EXISTS decisions (
id TEXT PRIMARY KEY,
type TEXT NOT NULL CHECK(type IN ('confirmed','question','rejected')),
domain TEXT NOT NULL,
title TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active'
CHECK(status IN ('active','superseded','resolved','open')),
date TEXT,
file_path TEXT NOT NULL,
synced_at TEXT NOT NULL DEFAULT (datetime('now')),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS decision_refs (
source_id TEXT NOT NULL REFERENCES decisions(id) ON DELETE CASCADE,
target_id TEXT NOT NULL REFERENCES decisions(id) ON DELETE CASCADE,
ref_type TEXT NOT NULL
CHECK(ref_type IN ('supersedes','references','resolves','depends_on','amends')),
note TEXT,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (source_id, target_id, ref_type)
);
-- Identity split (D-26): a ticket's stable, collision-proof identity is its
-- record_id (a locally-generated ULID, planning.NewRecordID); the friendly
-- T-NNN label lives in ticket_idmap and may be reconciled. Every structural
-- reference (parent, deps, history, labels) targets record_id, so a label
-- clash never corrupts the graph — only ticket_idmap needs a relabel.
CREATE TABLE IF NOT EXISTS tickets (
record_id TEXT PRIMARY KEY,
type TEXT NOT NULL CHECK(type IN ('initiative','epic','story','task','bug')),
parent_record_id TEXT REFERENCES tickets(record_id),
title TEXT NOT NULL,
description TEXT,
-- No CHECK enumeration: the ticket status vocabulary is per-vault
-- configurable (ticket_statuses in .pql/config.yaml). Validation lives
-- in Go (planning.StatusSet), so adding/renaming statuses needs no
-- schema change. The DEFAULT is a harmless fallback — CreateTicket
-- always inserts the configured default explicitly.
status TEXT NOT NULL DEFAULT 'backlog',
priority TEXT DEFAULT 'medium'
CHECK(priority IN ('critical','high','medium','low')),
assigned_to TEXT,
team TEXT,
decision_ref TEXT REFERENCES decisions(id),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
-- ticket_idmap maps a record_id to its current friendly label (T-NNN).
-- ticket_id is intentionally NOT globally unique: two uncoordinated clones
-- can mint the same label, which surfaces as a duplicate-label collision
-- (detected at replay) and is fixed with "pql ticket relabel".
CREATE TABLE IF NOT EXISTS ticket_idmap (
record_id TEXT PRIMARY KEY REFERENCES tickets(record_id),
ticket_id TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS ticket_deps (
blocker_record_id TEXT NOT NULL REFERENCES tickets(record_id),
blocked_record_id TEXT NOT NULL REFERENCES tickets(record_id),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (blocker_record_id, blocked_record_id)
);
CREATE TABLE IF NOT EXISTS ticket_history (
ticket_record_id TEXT NOT NULL REFERENCES tickets(record_id),
field TEXT NOT NULL,
old_value TEXT,
new_value TEXT,
changed_by TEXT,
changed_at TEXT NOT NULL DEFAULT (datetime('now')),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT UNIQUE,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS ticket_labels (
ticket_record_id TEXT NOT NULL REFERENCES tickets(record_id),
label TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (ticket_record_id, label)
);
CREATE TABLE IF NOT EXISTS meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL,
updated_at TEXT NOT NULL DEFAULT (datetime('now'))
);
CREATE INDEX IF NOT EXISTS idx_tickets_status ON tickets(status);
CREATE INDEX IF NOT EXISTS idx_tickets_team ON tickets(team);
CREATE INDEX IF NOT EXISTS idx_tickets_decision_ref ON tickets(decision_ref);
CREATE INDEX IF NOT EXISTS idx_tickets_assigned ON tickets(assigned_to);
CREATE INDEX IF NOT EXISTS idx_tickets_parent ON tickets(parent_record_id);
CREATE INDEX IF NOT EXISTS idx_ticket_idmap_label ON ticket_idmap(ticket_id);
CREATE INDEX IF NOT EXISTS idx_decisions_domain ON decisions(domain);
CREATE INDEX IF NOT EXISTS idx_decisions_type ON decisions(type);
CREATE INDEX IF NOT EXISTS idx_decision_refs_target ON decision_refs(target_id);
@@ -0,0 +1,139 @@
-- Auto-generated by pql init. CREATE TABLE statements
-- for the planning schema; per-table dir keeps the changelog
-- self-describing per D-15. CREATE TABLE IF NOT EXISTS is
-- idempotent so running schema files from each directory in
-- replay order is harmless.
--
-- Importer parses the markers below to detect schema drift
-- between the producing pql version and the local one — a
-- bumped canonical_version means projection rules changed
-- and replay must refuse rather than silently corrupt state.
-- pql:created_by: 2.2.0
-- pql:canonical_version: 2
CREATE TABLE IF NOT EXISTS decisions (
id TEXT PRIMARY KEY,
type TEXT NOT NULL CHECK(type IN ('confirmed','question','rejected')),
domain TEXT NOT NULL,
title TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active'
CHECK(status IN ('active','superseded','resolved','open')),
date TEXT,
file_path TEXT NOT NULL,
synced_at TEXT NOT NULL DEFAULT (datetime('now')),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS decision_refs (
source_id TEXT NOT NULL REFERENCES decisions(id) ON DELETE CASCADE,
target_id TEXT NOT NULL REFERENCES decisions(id) ON DELETE CASCADE,
ref_type TEXT NOT NULL
CHECK(ref_type IN ('supersedes','references','resolves','depends_on','amends')),
note TEXT,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (source_id, target_id, ref_type)
);
-- Identity split (D-26): a ticket's stable, collision-proof identity is its
-- record_id (a locally-generated ULID, planning.NewRecordID); the friendly
-- T-NNN label lives in ticket_idmap and may be reconciled. Every structural
-- reference (parent, deps, history, labels) targets record_id, so a label
-- clash never corrupts the graph — only ticket_idmap needs a relabel.
CREATE TABLE IF NOT EXISTS tickets (
record_id TEXT PRIMARY KEY,
type TEXT NOT NULL CHECK(type IN ('initiative','epic','story','task','bug')),
parent_record_id TEXT REFERENCES tickets(record_id),
title TEXT NOT NULL,
description TEXT,
-- No CHECK enumeration: the ticket status vocabulary is per-vault
-- configurable (ticket_statuses in .pql/config.yaml). Validation lives
-- in Go (planning.StatusSet), so adding/renaming statuses needs no
-- schema change. The DEFAULT is a harmless fallback — CreateTicket
-- always inserts the configured default explicitly.
status TEXT NOT NULL DEFAULT 'backlog',
priority TEXT DEFAULT 'medium'
CHECK(priority IN ('critical','high','medium','low')),
assigned_to TEXT,
team TEXT,
decision_ref TEXT REFERENCES decisions(id),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
-- ticket_idmap maps a record_id to its current friendly label (T-NNN).
-- ticket_id is intentionally NOT globally unique: two uncoordinated clones
-- can mint the same label, which surfaces as a duplicate-label collision
-- (detected at replay) and is fixed with "pql ticket relabel".
CREATE TABLE IF NOT EXISTS ticket_idmap (
record_id TEXT PRIMARY KEY REFERENCES tickets(record_id),
ticket_id TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS ticket_deps (
blocker_record_id TEXT NOT NULL REFERENCES tickets(record_id),
blocked_record_id TEXT NOT NULL REFERENCES tickets(record_id),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (blocker_record_id, blocked_record_id)
);
CREATE TABLE IF NOT EXISTS ticket_history (
ticket_record_id TEXT NOT NULL REFERENCES tickets(record_id),
field TEXT NOT NULL,
old_value TEXT,
new_value TEXT,
changed_by TEXT,
changed_at TEXT NOT NULL DEFAULT (datetime('now')),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT UNIQUE,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS ticket_labels (
ticket_record_id TEXT NOT NULL REFERENCES tickets(record_id),
label TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (ticket_record_id, label)
);
CREATE TABLE IF NOT EXISTS meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL,
updated_at TEXT NOT NULL DEFAULT (datetime('now'))
);
CREATE INDEX IF NOT EXISTS idx_tickets_status ON tickets(status);
CREATE INDEX IF NOT EXISTS idx_tickets_team ON tickets(team);
CREATE INDEX IF NOT EXISTS idx_tickets_decision_ref ON tickets(decision_ref);
CREATE INDEX IF NOT EXISTS idx_tickets_assigned ON tickets(assigned_to);
CREATE INDEX IF NOT EXISTS idx_tickets_parent ON tickets(parent_record_id);
CREATE INDEX IF NOT EXISTS idx_ticket_idmap_label ON ticket_idmap(ticket_id);
CREATE INDEX IF NOT EXISTS idx_decisions_domain ON decisions(domain);
CREATE INDEX IF NOT EXISTS idx_decisions_type ON decisions(type);
CREATE INDEX IF NOT EXISTS idx_decision_refs_target ON decision_refs(target_id);
+139
View File
@@ -0,0 +1,139 @@
-- Auto-generated by pql init. CREATE TABLE statements
-- for the planning schema; per-table dir keeps the changelog
-- self-describing per D-15. CREATE TABLE IF NOT EXISTS is
-- idempotent so running schema files from each directory in
-- replay order is harmless.
--
-- Importer parses the markers below to detect schema drift
-- between the producing pql version and the local one — a
-- bumped canonical_version means projection rules changed
-- and replay must refuse rather than silently corrupt state.
-- pql:created_by: 2.2.0
-- pql:canonical_version: 2
CREATE TABLE IF NOT EXISTS decisions (
id TEXT PRIMARY KEY,
type TEXT NOT NULL CHECK(type IN ('confirmed','question','rejected')),
domain TEXT NOT NULL,
title TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active'
CHECK(status IN ('active','superseded','resolved','open')),
date TEXT,
file_path TEXT NOT NULL,
synced_at TEXT NOT NULL DEFAULT (datetime('now')),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS decision_refs (
source_id TEXT NOT NULL REFERENCES decisions(id) ON DELETE CASCADE,
target_id TEXT NOT NULL REFERENCES decisions(id) ON DELETE CASCADE,
ref_type TEXT NOT NULL
CHECK(ref_type IN ('supersedes','references','resolves','depends_on','amends')),
note TEXT,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (source_id, target_id, ref_type)
);
-- Identity split (D-26): a ticket's stable, collision-proof identity is its
-- record_id (a locally-generated ULID, planning.NewRecordID); the friendly
-- T-NNN label lives in ticket_idmap and may be reconciled. Every structural
-- reference (parent, deps, history, labels) targets record_id, so a label
-- clash never corrupts the graph — only ticket_idmap needs a relabel.
CREATE TABLE IF NOT EXISTS tickets (
record_id TEXT PRIMARY KEY,
type TEXT NOT NULL CHECK(type IN ('initiative','epic','story','task','bug')),
parent_record_id TEXT REFERENCES tickets(record_id),
title TEXT NOT NULL,
description TEXT,
-- No CHECK enumeration: the ticket status vocabulary is per-vault
-- configurable (ticket_statuses in .pql/config.yaml). Validation lives
-- in Go (planning.StatusSet), so adding/renaming statuses needs no
-- schema change. The DEFAULT is a harmless fallback — CreateTicket
-- always inserts the configured default explicitly.
status TEXT NOT NULL DEFAULT 'backlog',
priority TEXT DEFAULT 'medium'
CHECK(priority IN ('critical','high','medium','low')),
assigned_to TEXT,
team TEXT,
decision_ref TEXT REFERENCES decisions(id),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
-- ticket_idmap maps a record_id to its current friendly label (T-NNN).
-- ticket_id is intentionally NOT globally unique: two uncoordinated clones
-- can mint the same label, which surfaces as a duplicate-label collision
-- (detected at replay) and is fixed with "pql ticket relabel".
CREATE TABLE IF NOT EXISTS ticket_idmap (
record_id TEXT PRIMARY KEY REFERENCES tickets(record_id),
ticket_id TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS ticket_deps (
blocker_record_id TEXT NOT NULL REFERENCES tickets(record_id),
blocked_record_id TEXT NOT NULL REFERENCES tickets(record_id),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (blocker_record_id, blocked_record_id)
);
CREATE TABLE IF NOT EXISTS ticket_history (
ticket_record_id TEXT NOT NULL REFERENCES tickets(record_id),
field TEXT NOT NULL,
old_value TEXT,
new_value TEXT,
changed_by TEXT,
changed_at TEXT NOT NULL DEFAULT (datetime('now')),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT UNIQUE,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS ticket_labels (
ticket_record_id TEXT NOT NULL REFERENCES tickets(record_id),
label TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (ticket_record_id, label)
);
CREATE TABLE IF NOT EXISTS meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL,
updated_at TEXT NOT NULL DEFAULT (datetime('now'))
);
CREATE INDEX IF NOT EXISTS idx_tickets_status ON tickets(status);
CREATE INDEX IF NOT EXISTS idx_tickets_team ON tickets(team);
CREATE INDEX IF NOT EXISTS idx_tickets_decision_ref ON tickets(decision_ref);
CREATE INDEX IF NOT EXISTS idx_tickets_assigned ON tickets(assigned_to);
CREATE INDEX IF NOT EXISTS idx_tickets_parent ON tickets(parent_record_id);
CREATE INDEX IF NOT EXISTS idx_ticket_idmap_label ON ticket_idmap(ticket_id);
CREATE INDEX IF NOT EXISTS idx_decisions_domain ON decisions(domain);
CREATE INDEX IF NOT EXISTS idx_decisions_type ON decisions(type);
CREATE INDEX IF NOT EXISTS idx_decision_refs_target ON decision_refs(target_id);
@@ -0,0 +1,139 @@
-- Auto-generated by pql init. CREATE TABLE statements
-- for the planning schema; per-table dir keeps the changelog
-- self-describing per D-15. CREATE TABLE IF NOT EXISTS is
-- idempotent so running schema files from each directory in
-- replay order is harmless.
--
-- Importer parses the markers below to detect schema drift
-- between the producing pql version and the local one — a
-- bumped canonical_version means projection rules changed
-- and replay must refuse rather than silently corrupt state.
-- pql:created_by: 2.2.0
-- pql:canonical_version: 2
CREATE TABLE IF NOT EXISTS decisions (
id TEXT PRIMARY KEY,
type TEXT NOT NULL CHECK(type IN ('confirmed','question','rejected')),
domain TEXT NOT NULL,
title TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active'
CHECK(status IN ('active','superseded','resolved','open')),
date TEXT,
file_path TEXT NOT NULL,
synced_at TEXT NOT NULL DEFAULT (datetime('now')),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS decision_refs (
source_id TEXT NOT NULL REFERENCES decisions(id) ON DELETE CASCADE,
target_id TEXT NOT NULL REFERENCES decisions(id) ON DELETE CASCADE,
ref_type TEXT NOT NULL
CHECK(ref_type IN ('supersedes','references','resolves','depends_on','amends')),
note TEXT,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (source_id, target_id, ref_type)
);
-- Identity split (D-26): a ticket's stable, collision-proof identity is its
-- record_id (a locally-generated ULID, planning.NewRecordID); the friendly
-- T-NNN label lives in ticket_idmap and may be reconciled. Every structural
-- reference (parent, deps, history, labels) targets record_id, so a label
-- clash never corrupts the graph — only ticket_idmap needs a relabel.
CREATE TABLE IF NOT EXISTS tickets (
record_id TEXT PRIMARY KEY,
type TEXT NOT NULL CHECK(type IN ('initiative','epic','story','task','bug')),
parent_record_id TEXT REFERENCES tickets(record_id),
title TEXT NOT NULL,
description TEXT,
-- No CHECK enumeration: the ticket status vocabulary is per-vault
-- configurable (ticket_statuses in .pql/config.yaml). Validation lives
-- in Go (planning.StatusSet), so adding/renaming statuses needs no
-- schema change. The DEFAULT is a harmless fallback — CreateTicket
-- always inserts the configured default explicitly.
status TEXT NOT NULL DEFAULT 'backlog',
priority TEXT DEFAULT 'medium'
CHECK(priority IN ('critical','high','medium','low')),
assigned_to TEXT,
team TEXT,
decision_ref TEXT REFERENCES decisions(id),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
-- ticket_idmap maps a record_id to its current friendly label (T-NNN).
-- ticket_id is intentionally NOT globally unique: two uncoordinated clones
-- can mint the same label, which surfaces as a duplicate-label collision
-- (detected at replay) and is fixed with "pql ticket relabel".
CREATE TABLE IF NOT EXISTS ticket_idmap (
record_id TEXT PRIMARY KEY REFERENCES tickets(record_id),
ticket_id TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS ticket_deps (
blocker_record_id TEXT NOT NULL REFERENCES tickets(record_id),
blocked_record_id TEXT NOT NULL REFERENCES tickets(record_id),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (blocker_record_id, blocked_record_id)
);
CREATE TABLE IF NOT EXISTS ticket_history (
ticket_record_id TEXT NOT NULL REFERENCES tickets(record_id),
field TEXT NOT NULL,
old_value TEXT,
new_value TEXT,
changed_by TEXT,
changed_at TEXT NOT NULL DEFAULT (datetime('now')),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT UNIQUE,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS ticket_labels (
ticket_record_id TEXT NOT NULL REFERENCES tickets(record_id),
label TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (ticket_record_id, label)
);
CREATE TABLE IF NOT EXISTS meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL,
updated_at TEXT NOT NULL DEFAULT (datetime('now'))
);
CREATE INDEX IF NOT EXISTS idx_tickets_status ON tickets(status);
CREATE INDEX IF NOT EXISTS idx_tickets_team ON tickets(team);
CREATE INDEX IF NOT EXISTS idx_tickets_decision_ref ON tickets(decision_ref);
CREATE INDEX IF NOT EXISTS idx_tickets_assigned ON tickets(assigned_to);
CREATE INDEX IF NOT EXISTS idx_tickets_parent ON tickets(parent_record_id);
CREATE INDEX IF NOT EXISTS idx_ticket_idmap_label ON ticket_idmap(ticket_id);
CREATE INDEX IF NOT EXISTS idx_decisions_domain ON decisions(domain);
CREATE INDEX IF NOT EXISTS idx_decisions_type ON decisions(type);
CREATE INDEX IF NOT EXISTS idx_decision_refs_target ON decision_refs(target_id);
+139
View File
@@ -0,0 +1,139 @@
-- Auto-generated by pql init. CREATE TABLE statements
-- for the planning schema; per-table dir keeps the changelog
-- self-describing per D-15. CREATE TABLE IF NOT EXISTS is
-- idempotent so running schema files from each directory in
-- replay order is harmless.
--
-- Importer parses the markers below to detect schema drift
-- between the producing pql version and the local one — a
-- bumped canonical_version means projection rules changed
-- and replay must refuse rather than silently corrupt state.
-- pql:created_by: 2.2.0
-- pql:canonical_version: 2
CREATE TABLE IF NOT EXISTS decisions (
id TEXT PRIMARY KEY,
type TEXT NOT NULL CHECK(type IN ('confirmed','question','rejected')),
domain TEXT NOT NULL,
title TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active'
CHECK(status IN ('active','superseded','resolved','open')),
date TEXT,
file_path TEXT NOT NULL,
synced_at TEXT NOT NULL DEFAULT (datetime('now')),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS decision_refs (
source_id TEXT NOT NULL REFERENCES decisions(id) ON DELETE CASCADE,
target_id TEXT NOT NULL REFERENCES decisions(id) ON DELETE CASCADE,
ref_type TEXT NOT NULL
CHECK(ref_type IN ('supersedes','references','resolves','depends_on','amends')),
note TEXT,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (source_id, target_id, ref_type)
);
-- Identity split (D-26): a ticket's stable, collision-proof identity is its
-- record_id (a locally-generated ULID, planning.NewRecordID); the friendly
-- T-NNN label lives in ticket_idmap and may be reconciled. Every structural
-- reference (parent, deps, history, labels) targets record_id, so a label
-- clash never corrupts the graph — only ticket_idmap needs a relabel.
CREATE TABLE IF NOT EXISTS tickets (
record_id TEXT PRIMARY KEY,
type TEXT NOT NULL CHECK(type IN ('initiative','epic','story','task','bug')),
parent_record_id TEXT REFERENCES tickets(record_id),
title TEXT NOT NULL,
description TEXT,
-- No CHECK enumeration: the ticket status vocabulary is per-vault
-- configurable (ticket_statuses in .pql/config.yaml). Validation lives
-- in Go (planning.StatusSet), so adding/renaming statuses needs no
-- schema change. The DEFAULT is a harmless fallback — CreateTicket
-- always inserts the configured default explicitly.
status TEXT NOT NULL DEFAULT 'backlog',
priority TEXT DEFAULT 'medium'
CHECK(priority IN ('critical','high','medium','low')),
assigned_to TEXT,
team TEXT,
decision_ref TEXT REFERENCES decisions(id),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
-- ticket_idmap maps a record_id to its current friendly label (T-NNN).
-- ticket_id is intentionally NOT globally unique: two uncoordinated clones
-- can mint the same label, which surfaces as a duplicate-label collision
-- (detected at replay) and is fixed with "pql ticket relabel".
CREATE TABLE IF NOT EXISTS ticket_idmap (
record_id TEXT PRIMARY KEY REFERENCES tickets(record_id),
ticket_id TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS ticket_deps (
blocker_record_id TEXT NOT NULL REFERENCES tickets(record_id),
blocked_record_id TEXT NOT NULL REFERENCES tickets(record_id),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (blocker_record_id, blocked_record_id)
);
CREATE TABLE IF NOT EXISTS ticket_history (
ticket_record_id TEXT NOT NULL REFERENCES tickets(record_id),
field TEXT NOT NULL,
old_value TEXT,
new_value TEXT,
changed_by TEXT,
changed_at TEXT NOT NULL DEFAULT (datetime('now')),
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT UNIQUE,
canonical_version INTEGER
);
CREATE TABLE IF NOT EXISTS ticket_labels (
ticket_record_id TEXT NOT NULL REFERENCES tickets(record_id),
label TEXT NOT NULL,
created_at TEXT NOT NULL DEFAULT (datetime('now')),
updated_at TEXT NOT NULL DEFAULT (datetime('now')),
deleted_at TEXT,
hash TEXT,
canonical_version INTEGER,
PRIMARY KEY (ticket_record_id, label)
);
CREATE TABLE IF NOT EXISTS meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL,
updated_at TEXT NOT NULL DEFAULT (datetime('now'))
);
CREATE INDEX IF NOT EXISTS idx_tickets_status ON tickets(status);
CREATE INDEX IF NOT EXISTS idx_tickets_team ON tickets(team);
CREATE INDEX IF NOT EXISTS idx_tickets_decision_ref ON tickets(decision_ref);
CREATE INDEX IF NOT EXISTS idx_tickets_assigned ON tickets(assigned_to);
CREATE INDEX IF NOT EXISTS idx_tickets_parent ON tickets(parent_record_id);
CREATE INDEX IF NOT EXISTS idx_ticket_idmap_label ON ticket_idmap(ticket_id);
CREATE INDEX IF NOT EXISTS idx_decisions_domain ON decisions(domain);
CREATE INDEX IF NOT EXISTS idx_decisions_type ON decisions(type);
CREATE INDEX IF NOT EXISTS idx_decision_refs_target ON decision_refs(target_id);
+92 -3
View File
@@ -4,6 +4,97 @@ 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.5.1] - 2026-08-11
### Fixed
- Backups no longer freeze the API. The config backup ran its ~21 minutes of tarring on the
event loop, so the whole service was unreachable 03:05-03:25 nightly; it now runs in a
worker thread. Health reports moved off the loop too.
- The health report's permission diagnostic prints the database's own message instead of
asserting a cause. It claimed the user lacked INSERT on `check_history` when that grant was
present and the missing one was USAGE on the sequence behind its serial id.
### Notes
- Requires `GRANT USAGE ON SEQUENCE check_history_id_seq TO scheduler_user`, applied
2026-08-11. The table grant alone does not permit the insert.
## [1.5.0] - 2026-08-11
### Added
- Backup executors report their own outcome to the homelab health record — one row in
`check_history` per run, success or failure. Replaces a monitor that inferred backup health
from file age and could not tell a failed backup from one that had not run yet.
### Notes
- Requires `GRANT INSERT ON check_history TO scheduler_user` in the `sysmon` database, applied
2026-08-11. Without it the report is refused, logged, and skipped; the backup itself is unaffected.
## [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
@@ -203,9 +294,7 @@ config_backup_executor.py 50% 📈
TOTAL 80% 🎯
```
## [Unreleased]
### Planned
## Planned
- Redis integration for distributed locking
- Webhook notifications for task completion
- Task dependencies (run task B after task A succeeds)
+160
View File
@@ -0,0 +1,160 @@
# CLAUDE.md — scheduler
"The Scheduler" — system-wide maintenance orchestration for tower-of-joy: config backups, doc
mirroring to Gitea, cleanup and retention, and arbitrary REST calls on a cron. Python 3.12 /
FastAPI, APScheduler, PostgreSQL and Redis. Container `scheduler` on `docker-dataplane`,
port **8090**, Redis DB **3**, Postgres DB **scheduler**.
It is the homelab's cron. Recurring work belongs here rather than in a systemd timer (workspace D-5).
## Live contract
`http://localhost:8090/openapi.json` — 10 paths, `version: 1.4.0` (verified 2026-08-09). Human
docs at `/docs`. Generated from running code, so read it instead of inferring routes.
**Every route except `/health` requires `Authorization: Bearer $SCHEDULER_API_KEY`.** An
unauthenticated call returns `{"detail": "Missing API key"}` at 200-shape JSON, not a 401 body
you might pattern-match on.
Note the old AGENTS.md told you to verify deploys against `192.168.86.149:8000/health`. That is
**tatlock's** port, not this service's. It is 8090.
## The thing that will mislead you: executors are chosen by data, not code
`src/executors/*.py` are **never statically imported**. `src/tasks/executor.py:202` does:
```python
module_path = f"src.executors.{executor_name}"
module = __import__(module_path, fromlist=['execute'])
```
where `executor_name` comes from a **row in the `scheduled_tasks` table**. Consequences, and
they defeat both of the usual checks:
- **grep finds nothing.** No file imports `config_backup_executor`; the name only ever exists as
a database string.
- **`sys.modules` finds nothing either.** A cold `import src.main` loads only `src`, `src.config`,
`src.main`, `src.models`, `src.tasks`, `src.tasks.executor`. Every executor is absent until a
task actually fires. Absence there is a timing artifact, not evidence of death.
**The authoritative source is the database.** As of 2026-08-09:
| Executor | Rows | Enabled |
|---|---|---|
| `rest_api_executor` | 16 | yes |
| `doc_sync_executor` | 2 | yes |
| `config_backup_executor`, `docker_prune_executor`, `gitea_release_cleanup_executor`, `portainer_backup_executor`, `postgres_retention_executor` | 1 each | yes |
| `example_executor` | 1 | **no** |
| `gcs_backup_executor` | **0** | — |
`gcs_backup_executor` has no rows at all. That does **not** make it dead code: it becomes live
the instant someone inserts a row naming it, with no code change and no deploy. Treat unreferenced
executors as *dormant*, not removable. An executor's contract is a module-level
`execute(config, settings)` — a missing one is caught at run time and reported as
`Executor <name> missing execute() function`, not at import or startup.
Re-check with the query rather than trusting the table above:
```bash
docker exec scheduler python3 -c "
import psycopg2
from src.config import get_settings
s = get_settings()
c = psycopg2.connect(host=s.postgres_host, port=s.postgres_port, dbname=s.postgres_db,
user=s.postgres_user, password=s.postgres_password)
cur = c.cursor()
cur.execute('SELECT executor, count(*), bool_or(enabled) FROM scheduled_tasks GROUP BY executor ORDER BY 1')
[print(r) for r in cur.fetchall()]"
```
Build the connection from `get_settings()` fields as above. Do not print the assembled URL — it
carries the Postgres password.
## Database
Three tables, and the names do not match the API paths: **`scheduled_tasks`** (not `tasks`
`SELECT … FROM tasks` fails with `UndefinedTable`), `task_executions`, `doc_sources`. Schema is
SQLAlchemy (`src/models.py`); there is no Alembic here, unlike core-api.
## Layout, and one trap in it
`src/main.py` (app + routes), `src/config.py` (pydantic-settings), `src/models.py`,
`src/tasks/executor.py` (the scheduling engine), `src/executors/` (the dynamically-loaded units).
**`src/config/` also exists and is an empty directory.** `import src.config` resolves to
`src/config.py` — verified in the container, `__file__` is `/app/src/config.py`, because a
regular module wins over a namespace package. Do not "fix" this by moving config into the
directory, and do not assume the directory is a package with contents.
## Registering tasks
Tasks are DB-driven, registered over the API — not YAML, not a file in this repo. See
`TASK_REGISTRATION.md` for the payload shape and the cron-field conventions (`hour: -1` means
every hour). There is also a workspace-level `scheduler` skill for driving it conversationally.
## Working here
Group new work by domain rather than by file type; a single large `routers/` folder is the thing
to avoid. Reference: [FastAPI best practices](https://github.com/zhanymkanov/fastapi-best-practices).
```bash
.venv/bin/python -m pytest tests/ # or: pytest tests/
```
Test dependencies are the `test` extra in `pyproject.toml` (pytest, pytest-asyncio, pytest-cov,
freezegun). `pytest.ini` is at the repo root. No linter is configured — no ruff/flake8 config and
neither in the dependencies — so do not assume `ruff check` exists here.
## CI
`.gitea/workflows/build.yml` is the only workflow and triggers **only on `v*` tag push**: build,
push image, ping Watchtower. There is **no CI test or lint gate**. Run the tests yourself before
tagging.
## Work tracking
Work lives in **pql**, not a markdown TODO. **This repo's vault is standalone** — its tickets
and its internal decisions live here in `.pql/` and `governance/`, and travel with a clone,
because `.pql/changelog/` is committed and replayed by the git hooks (workspace D-15). The databases are
gitignored and rebuildable with `pql plan rebuild`.
`pql` is **not** on the non-interactive `PATH` — invoke it as `/home/jpmschweitzer/.local/bin/pql`.
From inside this repo no `--vault` is needed: pql anchors at the nearest `.git/` ancestor, which
is this repo.
```bash
/home/jpmschweitzer/.local/bin/pql ticket list # this repo's open work
/home/jpmschweitzer/.local/bin/pql plan whatsnext # next unblocked item, with context
/home/jpmschweitzer/.local/bin/pql decisions list # this repo's own decisions
```
Stack-level decisions that constrain this service — the host, the network, deploy mechanics,
and the fact that recurring work belongs here at all (workspace D-5) — live in the **workspace** vault
and need the flag:
```bash
/home/jpmschweitzer/.local/bin/pql --vault /mnt/media/Projects decisions list --domain scheduler
```
Note `ticket new --decision D-N` resolves ids within **one** vault, so a ticket here cannot link
to a workspace decision. Cite the id in the ticket body instead.
## Git
- **History is linear — no merge commits.** Work on `main`, or a short-lived branch that is
fast-forwarded and deleted. The previous AGENTS.md mandated a feature branch per change; that
rule was retired workspace-wide on 2026-08-08.
- **Conventional Commits**: `feat:`, `fix:`, `refactor:`, `docs:`, `chore:`.
- **Stage explicitly. Never `git add -A`** — it is denied by policy, and it sweeps in whatever
else is dirty, including secrets.
- Update `CHANGELOG.md` for every user-facing change, under `[Unreleased]`.
## Releasing
Ask whether a deploy is wanted first — it is not automatic.
1. Bump `version` in `pyproject.toml` (patch for fixes, minor for features).
2. Move `[Unreleased]` entries into a dated section in `CHANGELOG.md`.
3. Stage the changed files by name, commit, tag `vX.Y.Z`, `git push origin main --tags`.
4. Gitea CI builds and pushes on the tag; Watchtower deploys it.
5. Verify: `curl http://192.168.86.149:8090/health`.
+55
View File
@@ -0,0 +1,55 @@
# scheduler — the repo's command surface (D-27).
#
# There is no venv in this working tree today, even though CLAUDE.md documents
# `.venv/bin/python -m pytest`. `make test` says so rather than failing with a
# bare "No such file or directory", and `make setup` creates one.
#
# `python3` on this host is 3.8; PYTHON names 3.12 explicitly (D-26).
VENV := $(CURDIR)/.venv
PYTHON ?= python3.12
.DEFAULT_GOAL := help
.PHONY: help
help: ## Show this help
@grep -hE '^[a-z][a-z0-9_-]*:.*?## ' $(MAKEFILE_LIST) \
| awk 'BEGIN{FS=":.*?## "}{printf " \033[36m%-14s\033[0m %s\n", $$1, $$2}'
.PHONY: setup
setup: ## Create the venv and install the test extra
$(PYTHON) -m venv .venv
$(VENV)/bin/pip install -e ".[test]"
.PHONY: test
test: ## Run the test suite
@test -x $(VENV)/bin/python || { echo "FAIL — no venv in this tree; run: make setup"; exit 69; }
$(VENV)/bin/python -m pytest tests/
# No `lint` target, deliberately. CLAUDE.md states it outright: no linter is
# configured, no ruff or flake8 config, neither in the dependencies. Per D-27
# the name is reserved for repos that lint rather than mandated everywhere — a
# target here could only fail or report clean for something never run.
# git hands a hook a non-login shell, which never sees ~/.local/bin — where
# gitleaks lands. Without this the scan reports "not installed" on every push,
# which is a check that fails open (D-24).
export PATH := $(HOME)/.local/bin:/usr/local/bin:$(PATH)
.PHONY: secrets
secrets: ## Scan the commits about to be pushed for credentials
@ci/secrets.sh
# The call surface is identical in every repo; what it runs is not.
#
# `secrets` runs first, deliberately: it is the only failure here that cannot be
# undone by fixing it afterwards. A failed lint costs another commit; a pushed
# credential is cached and indexed whether or not it is later deleted.
#
# Some of these fail today, and are left wired anyway. The state was measured
# once and written down in T-56 rather than being worked around here — a gate
# quietly narrowed to what already passes is a gate that reports success for
# doing nothing, which is the failure this workspace keeps rediscovering.
.PHONY: pre-push
pre-push: secrets ## Everything the pre-push hook runs
@echo " -- not gated here yet: lint (no linter configured) and test (T-56)"
+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
Executable
+50
View File
@@ -0,0 +1,50 @@
#!/usr/bin/env bash
# Secret scan over the commits about to be pushed.
#
# Lives here rather than inside .githooks/pre-push so it can be read, run by
# hand (`make secrets`), and changed under review. A hook is a trigger; it is
# not a home for logic. Identical in every repo in this workspace (D-27).
set -euo pipefail
cd "$(git rev-parse --show-toplevel)"
# A non-login shell — which is what git gives a hook — skips /etc/profile.d
# and never sees ~/.local/bin, where the gitleaks release tarball lands.
# Without this the scan reports "not installed" on every push.
[ -d "$HOME/.local/bin" ] && PATH="$HOME/.local/bin:$PATH"
if ! command -v gitleaks >/dev/null 2>&1; then
echo "FAIL secrets — gitleaks not installed, so this check would be a no-op pretending to pass." >&2
echo " https://github.com/gitleaks/gitleaks/releases → ~/.local/bin/gitleaks" >&2
exit 1
fi
# Scan the outgoing range, not full history. History here carries findings
# that are settled — test fixtures and vendored third-party code — and a gate
# that fails on something unfixable gets bypassed within a week. What matters
# is what is about to leave this machine.
if upstream=$(git rev-parse --abbrev-ref --symbolic-full-name '@{u}' 2>/dev/null); then
range="$upstream..HEAD"
elif git rev-parse --verify --quiet origin/main >/dev/null; then
range="origin/main..HEAD"
else
range=""
fi
if [ -z "$range" ]; then
gitleaks dir . --redact --no-banner --exit-code 1 || {
echo "FAIL secrets — gitleaks found a credential in the working tree." >&2; exit 1; }
exit 0
fi
[ -n "$(git log --oneline "$range" 2>/dev/null)" ] || exit 0
gitleaks git . --log-opts="$range" --redact --no-banner --exit-code 1 >/dev/null 2>&1 || {
echo "FAIL secrets — gitleaks found a credential in the commits being pushed." >&2
echo " inspect (values redacted): gitleaks git . --log-opts=\"$range\" --redact" >&2
echo " then remove and rotate it, or suppress deliberately:" >&2
echo " inline '# gitleaks:allow <reason>'" >&2
echo " or add the fingerprint to .gitleaksignore WITH a reason" >&2
exit 1
}
echo " ok secrets"
+54
View File
@@ -0,0 +1,54 @@
# Decisions, Questions, Rejected
This directory holds structured planning records that pql parses
into pql.db. Each record is a `### [DQR]-N: Title` heading inside
a markdown file. Files live in three per-type subdirectories:
- `decisions/<domain>.md` — confirmed design decisions
- `questions/<domain>.md` — open questions that may resolve into
decisions or rejected proposals
- `rejected/<domain>.md` — rejected proposals (kept for the audit
trail)
The parser infers domain from the filename stem and record type
from the parent subdirectory.
D-records that propose implementation work link to `initiative`-type
tickets via `decision_ref`. Run `pql decisions show <id>
--with-tickets` to inspect implementation status.
## Recommended domains
Start with this canonical set; create files as records land in
each domain:
- **architecture** — structural commitments (storage, layering,
languages, libraries)
- **process** — team workflow (commits, branches, releases, reviews)
- **design** — user-facing surface (UX, UI, public APIs)
- **coding-conventions** — team-internal code shape (style, lint,
file layout)
- **testing** — quality strategy (coverage, layers, gates)
You might also want, project-permitting:
- `accessibility` — if you ship user-facing software
- `security` — if you handle user data or network surfaces
- `licensing` — if you release open-source or commercial
- `documentation` — if user-docs are non-trivial
- `deployment` — if shipping is non-trivial
- `performance` — if you have perf budgets / SLOs
<!-- pql:records (auto-generated; do not edit manually) -->
## Decisions
- _(none)_
## Open questions
- _(none)_
## Rejected
- _(none)_
+3 -1
View File
@@ -1,6 +1,6 @@
[project]
name = "the-scheduler"
version = "1.0.4"
version = "1.5.1"
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 -53
View File
@@ -2,6 +2,20 @@
Config Backup Executor
Backs up Docker container configs and host-based service configs.
Replicates functionality of maintenance container's backup-configs.sh
The work runs in a worker thread. `execute` is awaited by the scheduling engine
on the same event loop that serves `/health` and the whole REST API, and this
job spends ~21 minutes inside tarfile and zlib. Called inline it starves the
loop for that entire window: on 2026-08-11 the API was unreachable 03:05-03:25
nightly and again at 07:39 when the job was triggered by hand, which is T-74.
The hourly health check fires at :35 and had never once sampled the outage.
`asyncio.to_thread` is the whole fix — zlib releases the GIL while compressing,
so the loop gets scheduled normally. One caveat worth knowing: the engine wraps
executors in `asyncio.wait_for`, and a thread cannot be cancelled. On timeout
the task is recorded as failed while the tar keeps running to completion. That
is still strictly better than blocking everything, and the configured timeout
(3600s) is well clear of the observed 1263s.
"""
import asyncio
import logging
@@ -9,16 +23,75 @@ import tarfile
import tempfile
from datetime import datetime, timedelta
from pathlib import Path
from typing import List, Dict, Any
from typing import Any, Callable, Dict, List, Optional
from src.config import Settings
from src.executors import health_report
logger = logging.getLogger(__name__)
async def execute(config: dict, settings: Settings) -> str:
def _create_tar_filter(excludes: List[str]) -> Callable[[Any], Optional[Any]]:
"""Build the tarfile filter for a set of exclude patterns.
Matching is **substring**, not glob, and that is deliberate: a pattern is
reduced to its literal core by dropping leading `*/` and trailing `/*`, and
a member is excluded when that core appears anywhere in its name.
It looks like a half-finished glob and the temptation is to "fix" it with
`fnmatch`. Doing so would break the deployed configuration badly, because
that configuration is written against these semantics:
".log" fnmatch would match only a file named exactly `.log`,
so every log file starts being archived instead.
"ollama/models/*" fnmatch anchors at the start of the name, and members
are named `docker-data/ollama/models/...`, so nothing
matches and many GB of model blobs enter the archive.
Same for `amp/Versions/*`, `qdrant/storage/*` and the rest — every one is an
unanchored mid-path fragment. The nightly archive would grow, not shrink.
`tests/test_config_backup_executor.py` pins both cases so this cannot be
changed silently.
"""
Execute config backup task.
cores = [p.replace('*/', '').replace('/*', '') for p in excludes]
def tar_filter(tarinfo):
for core in cores:
if core in tarinfo.name:
logger.debug(f"Excluding: {tarinfo.name}")
return None
return tarinfo
return tar_filter
def _cleanup_old_backups(backup_dir: Path, retention_days: int) -> None:
"""Remove backups older than the retention period."""
cutoff_date = datetime.now() - timedelta(days=retention_days)
removed_count = 0
removed_size = 0
logger.info(f"Cleaning up backups older than {retention_days} days...")
for backup_file in backup_dir.glob('docker-configs-*.tar.gz'):
file_mtime = datetime.fromtimestamp(backup_file.stat().st_mtime)
if file_mtime < cutoff_date:
file_size = backup_file.stat().st_size
logger.info(f"Removing old backup: {backup_file.name} (from {file_mtime:%Y-%m-%d})")
backup_file.unlink()
removed_count += 1
removed_size += file_size
if removed_count > 0:
removed_size_mb = removed_size / (1024 * 1024)
logger.info(f"Removed {removed_count} old backups, freed {removed_size_mb:.2f}MB")
else:
logger.info("No old backups to remove")
def _backup(config: dict) -> str:
"""The blocking body of the backup. Runs in a worker thread, never on the loop.
Config schema:
{
@@ -34,20 +107,11 @@ async def execute(config: dict, settings: Settings) -> str:
"compress": true
}
Args:
config: Backup configuration
settings: Global scheduler settings
Returns:
Summary of backup operation
Raises:
Exception: On backup failure
Returns a one-line summary. Raises on failure.
"""
sources = config.get('sources', [])
backup_dir = Path(config.get('backup_dir', '/backups/docker-configs'))
retention_days = config.get('retention_days', 30)
compress = config.get('compress', True)
if not sources:
raise ValueError("No backup sources configured")
@@ -58,15 +122,12 @@ async def execute(config: dict, settings: Settings) -> str:
logger.info(f"Starting Docker configs backup: {backup_filename}")
# Create backup directory
backup_dir.mkdir(parents=True, exist_ok=True)
# Create temporary directory for staging
with tempfile.TemporaryDirectory(prefix='backup-') as temp_dir:
temp_path = Path(temp_dir)
results = []
# Backup each source
for source in sources:
source_path = Path(source['path'])
source_name = source['name']
@@ -78,23 +139,13 @@ async def execute(config: dict, settings: Settings) -> str:
logger.info(f"Backing up {source_name} from {source_path}")
# Create tar for this source
source_tar = temp_path / f"{source_name}.tar.gz"
def tar_filter(tarinfo):
"""Filter function to exclude patterns."""
for pattern in excludes:
# Simple pattern matching (could be enhanced with fnmatch)
if pattern.replace('*/', '').replace('/*', '') in tarinfo.name:
logger.debug(f"Excluding: {tarinfo.name}")
return None
return tarinfo
with tarfile.open(source_tar, 'w:gz') as tar:
tar.add(
source_path,
arcname=source_name,
filter=tar_filter,
filter=_create_tar_filter(excludes),
recursive=True
)
@@ -102,23 +153,19 @@ async def execute(config: dict, settings: Settings) -> str:
results.append(f"{source_name}: {source_size:.2f}MB")
logger.info(f"Backed up {source_name}: {source_size:.2f}MB")
# Combine all source backups into final archive
logger.info("Creating combined backup archive...")
with tarfile.open(backup_file, 'w:gz') as final_tar:
for item in temp_path.glob('*.tar.gz'):
final_tar.add(item, arcname=item.name)
# Verify backup created
if not backup_file.exists():
raise Exception("Backup file was not created")
backup_size = backup_file.stat().st_size / (1024 * 1024) # MB
logger.info(f"Backup created successfully: {backup_size:.2f}MB")
# Clean up old backups
await cleanup_old_backups(backup_dir, retention_days)
_cleanup_old_backups(backup_dir, retention_days)
# Count remaining backups
backup_count = len(list(backup_dir.glob('docker-configs-*.tar.gz')))
total_size = sum(f.stat().st_size for f in backup_dir.glob('docker-configs-*.tar.gz'))
total_size_mb = total_size / (1024 * 1024)
@@ -133,27 +180,36 @@ async def execute(config: dict, settings: Settings) -> str:
return output
async def cleanup_old_backups(backup_dir: Path, retention_days: int):
"""Remove backups older than retention period."""
cutoff_date = datetime.now() - timedelta(days=retention_days)
removed_count = 0
removed_size = 0
async def _run(config: dict, settings: Settings) -> str:
"""Await the backup without holding the event loop. See the module docstring."""
return await asyncio.to_thread(_backup, config)
logger.info(f"Cleaning up backups older than {retention_days} days...")
for backup_file in backup_dir.glob('docker-configs-*.tar.gz'):
# Get file modification time
file_mtime = datetime.fromtimestamp(backup_file.stat().st_mtime)
async def execute(config: dict, settings: Settings) -> str:
"""Run the backup and report its own outcome to check_history (D-33, T-69).
if file_mtime < cutoff_date:
file_size = backup_file.stat().st_size
logger.info(f"Removing old backup: {backup_file.name} (from {file_mtime:%Y-%m-%d})")
backup_file.unlink()
removed_count += 1
removed_size += file_size
if removed_count > 0:
removed_size_mb = removed_size / (1024 * 1024)
logger.info(f"Removed {removed_count} old backups, freed {removed_size_mb:.2f}MB")
else:
logger.info("No old backups to remove")
The report wraps the work rather than living inside it, so the failure path
cannot be forgotten: an exception is reported as critical and then re-raised,
leaving the task's own status untouched. Reporting only success would
reproduce exactly the blind spot this replaces — a monitor that cannot tell
a failed backup from one that has not run.
"""
try:
output = await _run(config, settings)
except Exception as exc:
await health_report.report_async(
settings,
domain="backup",
status=health_report.CRITICAL,
source="scheduler/config_backup_executor",
metrics={"job": "scheduler/config_backup_executor", "error": str(exc)[:400]},
)
raise
await health_report.report_async(
settings,
domain="backup",
status=health_report.OK,
source="scheduler/config_backup_executor",
metrics={"job": "scheduler/config_backup_executor", "summary": output[:400]},
)
return output
+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)
+155
View File
@@ -0,0 +1,155 @@
"""Report a task's own outcome to the homelab's central health record.
Why a process reports itself, rather than a monitor inferring it:
The sysmon `backup` domain used to poll the mtime of the newest file in the
backup directory, hourly, against a 48-hour threshold — for a job that runs once
a day. Forty-seven of every forty-eight runs could not produce a new answer, and
worse, a file-age poll cannot distinguish "the backup failed" from "the backup
has not run yet". If tonight's job dies, yesterday's archive is 24 hours old and
still reads healthy, and keeps reading healthy until hour 48. A failure stayed
invisible for two days to the check whose only job was noticing it.
This executor knows at 03:05. So it says so.
Recorded as D-33 in the workspace vault: `check_history` is the central health
record and any self-maintained service may push a row describing its own
outcome. sysmon polls only the things that cannot report themselves.
Three consequences that are load-bearing here:
- `source` names the producer, because the table now has several writers and a
row must say which one wrote it.
- `domain` is a shared namespace. Two producers claiming one name would
interleave silently.
- **A reporting failure must never fail the task.** Backing up successfully and
failing to mention it is strictly better than the reverse. Everything here is
caught and logged, which means the absence of rows is the only symptom a
broken reporter produces — so check for rows, not for errors.
"""
import asyncio
import json
import logging
from datetime import datetime, timezone
from typing import Any, Dict, Optional
import psycopg2
logger = logging.getLogger(__name__)
# The database holding check_history. Not the scheduler's own database — this is
# a cross-service write into the health record, and it is deliberate (D-33).
HEALTH_DB = "sysmon"
OK = "ok"
WARNING = "warning"
CRITICAL = "critical"
def report(
settings: Any,
domain: str,
status: str,
source: str,
metrics: Optional[Dict[str, Any]] = None,
) -> bool:
"""Write one row to check_history. Returns whether it landed.
Never raises. A caller that lets this failure surface would turn a
successful backup into a failed task, which inverts the point.
"""
metrics = metrics or {}
now = datetime.now(timezone.utc)
result = {
# The envelope the table has carried since the shell era. A reader of a
# year of history should not have to know which producer wrote a row in
# order to parse it.
"timestamp": now.isoformat(),
"source": source,
"domain": domain,
"status": status,
"metrics": metrics,
}
try:
conn = psycopg2.connect(
host=settings.postgres_host,
port=settings.postgres_port,
database=HEALTH_DB,
user=settings.postgres_user,
password=settings.postgres_password,
connect_timeout=10,
)
except Exception as exc: # noqa: BLE001 - reporting must not raise
logger.warning("health report for %s could not connect to %s: %s", domain, HEALTH_DB, exc)
return False
try:
with conn:
with conn.cursor() as cur:
# Unqualified table name, resolved through the search_path of the
# sysmon database. Qualifying it as sysmon.check_history looks
# more careful and is wrong — that schema does not exist.
cur.execute(
"INSERT INTO check_history (host, domain, status, ts, result) "
"VALUES (%s, %s, %s, %s, %s)",
(_host(), domain, status, now, json.dumps(result)),
)
logger.info("health report: %s=%s recorded", domain, status)
return True
except psycopg2.errors.InsufficientPrivilege as exc:
# Named separately because it is the expected first failure and the fix
# is a grant rather than a code change. The database's own message is
# printed verbatim rather than summarised: the first version of this
# asserted "lacks INSERT on check_history" and was wrong — the table
# grant was present and what was actually missing was USAGE on
# check_history_id_seq, the sequence behind its serial id. A diagnostic
# that names a cause it did not observe sends the reader to the wrong
# fix with confidence.
#
# GRANT INSERT ON check_history TO <user>;
# GRANT USAGE ON SEQUENCE check_history_id_seq TO <user>;
logger.warning(
"health report for %s refused by the database: %s. The task itself succeeded; "
"only the report was lost.",
domain, str(exc).strip().splitlines()[0],
)
return False
except Exception as exc: # noqa: BLE001 - reporting must not raise
logger.warning("health report for %s failed: %s", domain, exc)
return False
finally:
conn.close()
async def report_async(
settings: Any,
domain: str,
status: str,
source: str,
metrics: Optional[Dict[str, Any]] = None,
) -> bool:
"""`report` for callers on the event loop. Prefer this one inside executors.
psycopg2 is a blocking driver, so calling `report` directly from an
`async def` holds the loop for the length of the connect and insert — up to
`connect_timeout` seconds if the database is unreachable, which is exactly
when a report is most likely to be attempted. The scheduler serves its own
`/health` from that loop, so the cost of a slow report is the whole service
appearing down (T-74).
"""
return await asyncio.to_thread(report, settings, domain, status, source, metrics)
def _host() -> str:
"""The host a row is attributed to.
Every Redis key and check_history row is scoped by host so a second machine
reporting into the same store stays distinguishable. The scheduler runs in a
container, whose hostname is a container id — useless as an attribution — so
the physical host is named explicitly.
"""
import os
return os.environ.get("SYSMON_HOST", "tower-of-joy")
+171
View File
@@ -0,0 +1,171 @@
"""
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
from src.executors import health_report
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 _run(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
async def execute(config: dict, settings: Settings) -> str:
"""Run the backup and report its own outcome to check_history (D-33, T-69).
The report wraps the work rather than living inside it, so the failure path
cannot be forgotten: an exception is reported as critical and then re-raised,
leaving the task's own status untouched. Reporting only success would
reproduce exactly the blind spot this replaces — a monitor that cannot tell
a failed backup from one that has not run.
"""
try:
output = await _run(config, settings)
except Exception as exc:
await health_report.report_async(
settings,
domain="backup",
status=health_report.CRITICAL,
source="scheduler/portainer_backup_executor",
metrics={"job": "scheduler/portainer_backup_executor", "error": str(exc)[:400]},
)
raise
await health_report.report_async(
settings,
domain="backup",
status=health_report.OK,
source="scheduler/portainer_backup_executor",
metrics={"job": "scheduler/portainer_backup_executor", "summary": output[:400]},
)
return output
@@ -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
+207 -180
View File
@@ -1,214 +1,241 @@
"""Tests for the config backup executor.
These replace a set that arrived with the portainer-core extraction and had
never passed in this repo: they patched the `Path` class wholesale, asserted
`"backed up" in result` against a function that returns `"Backup completed: …"`,
and one wrapped its only call in `except Exception: pass` with its assertion
commented out. Seven were red from the initial commit onward, and there is no CI
test gate here to notice (see CLAUDE.md).
The replacements use real directories and real archives under `tmp_path`. A
backup executor's whole job is what ends up inside the tar, and mocking
`tarfile` means nothing is checked.
"""
Tests for the config backup executor.
"""
import pytest
import asyncio
import os
import tarfile
import time
from pathlib import Path
from unittest.mock import AsyncMock, patch, MagicMock, mock_open
from src.executors import config_backup_executor
import pytest
from src.config import Settings
from src.executors import config_backup_executor
def _make_source(root: Path) -> Path:
"""A source tree with one file per exclude pattern the deployment uses."""
src = root / "docker-data"
(src / "ollama" / "models" / "blobs").mkdir(parents=True)
(src / "svc" / "cache").mkdir(parents=True)
(src / "svc" / "logs").mkdir(parents=True)
(src / "keep.conf").write_text("keep me")
(src / "svc" / "settings.json").write_text("keep me too")
(src / "ollama" / "models" / "blobs" / "sha256-abc").write_text("many GB in reality")
(src / "svc" / "cache" / "junk.bin").write_text("disposable")
(src / "svc" / "logs" / "app.log").write_text("noisy")
return src
def _members(archive: Path) -> set:
"""Names inside the inner per-source tar of a combined backup archive."""
with tarfile.open(archive, "r:gz") as outer:
inner_name = outer.getnames()[0]
fh = outer.extractfile(inner_name)
with tarfile.open(fileobj=fh, mode="r:gz") as inner:
return set(inner.getnames())
@pytest.mark.executor
@pytest.mark.unit
class TestConfigBackupExecutor:
"""Tests for config_backup_executor module."""
class TestConfigBackup:
@pytest.mark.asyncio
async def test_executor_requires_config_fields(self, test_settings: Settings):
"""Test that executor validates required config."""
incomplete_config = {
"sources": []
# Missing backup_dir
}
async def test_no_sources_is_rejected(self, test_settings: Settings, tmp_path: Path):
with pytest.raises(ValueError, match="No backup sources"):
await config_backup_executor._run({"backup_dir": str(tmp_path)}, test_settings)
with pytest.raises((ValueError, KeyError)):
await config_backup_executor.execute(incomplete_config, test_settings)
def test_archive_contains_the_source(self, tmp_path: Path):
src = _make_source(tmp_path)
out = tmp_path / "backups"
summary = config_backup_executor._backup({
"sources": [{"path": str(src), "name": "docker-data", "excludes": []}],
"backup_dir": str(out),
})
@pytest.mark.asyncio
async def test_executor_with_minimal_config(self, test_settings: Settings, tmp_path: Path):
"""Test executor with minimal valid configuration."""
backup_dir = tmp_path / "backups"
backup_dir.mkdir()
assert "Backup completed" in summary
archives = list(out.glob("docker-configs-*.tar.gz"))
assert len(archives) == 1
assert "docker-data/keep.conf" in _members(archives[0])
source_dir = tmp_path / "source"
source_dir.mkdir()
(source_dir / "test.txt").write_text("test content")
config = {
def test_excludes_keep_matching_members_out(self, tmp_path: Path):
"""The patterns here are the ones the deployed task actually carries."""
src = _make_source(tmp_path)
out = tmp_path / "backups"
config_backup_executor._backup({
"sources": [{
"path": str(source_dir),
"name": "test-source",
"excludes": []
"path": str(src),
"name": "docker-data",
"excludes": ["*/cache/*", ".log", "ollama/models/*"],
}],
"backup_dir": str(backup_dir),
"compress": True,
"retention_days": 30
}
"backup_dir": str(out),
})
with patch('src.executors.config_backup_executor.Path') as mock_path_cls:
# Setup path mocking
mock_source = MagicMock()
mock_source.exists.return_value = True
mock_source.is_dir.return_value = True
mock_source.iterdir.return_value = [MagicMock(name="test.txt")]
names = _members(next(iter(out.glob("docker-configs-*.tar.gz"))))
assert "docker-data/keep.conf" in names
assert "docker-data/svc/settings.json" in names
assert "docker-data/svc/cache/junk.bin" not in names
assert "docker-data/svc/logs/app.log" not in names
assert "docker-data/ollama/models/blobs/sha256-abc" not in names
mock_backup = MagicMock()
mock_backup.mkdir = MagicMock()
def test_missing_source_is_skipped_not_fatal(self, tmp_path: Path):
src = _make_source(tmp_path)
out = tmp_path / "backups"
summary = config_backup_executor._backup({
"sources": [
{"path": str(tmp_path / "does-not-exist"), "name": "gone", "excludes": []},
{"path": str(src), "name": "docker-data", "excludes": []},
],
"backup_dir": str(out),
})
assert "docker-data" in summary
assert "gone" not in summary
def path_side_effect(p):
if str(p) == str(source_dir):
return mock_source
elif str(p) == str(backup_dir):
return mock_backup
return MagicMock()
def test_cleanup_removes_only_expired_backups(self, tmp_path: Path):
out = tmp_path / "backups"
out.mkdir()
old = out / "docker-configs-20200101-000000.tar.gz"
recent = out / "docker-configs-20991231-000000.tar.gz"
unrelated = out / "notes.txt"
for f in (old, recent, unrelated):
f.write_text("x")
mock_path_cls.side_effect = path_side_effect
long_ago = time.time() - (30 * 86400)
os.utime(old, (long_ago, long_ago))
with patch('tarfile.open'), \
patch('src.executors.config_backup_executor._cleanup_old_backups'):
config_backup_executor._cleanup_old_backups(out, retention_days=7)
result = await config_backup_executor.execute(config, test_settings)
assert "backed up" in result.lower() or "success" in result.lower()
@pytest.mark.asyncio
async def test_executor_excludes_patterns(self, test_settings: Settings, tmp_path: Path):
"""Test that executor respects exclude patterns."""
backup_dir = tmp_path / "backups"
source_dir = tmp_path / "source"
config = {
"sources": [{
"path": str(source_dir),
"name": "test",
"excludes": ["*.log", "cache/*"]
}],
"backup_dir": str(backup_dir),
"compress": True
}
with patch('src.executors.config_backup_executor.Path'), \
patch('tarfile.open') as mock_tar, \
patch('src.executors.config_backup_executor._cleanup_old_backups'):
# Mock tarfile
mock_tar_obj = MagicMock()
mock_tar.return_value.__enter__ = MagicMock(return_value=mock_tar_obj)
mock_tar.return_value.__exit__ = MagicMock(return_value=None)
try:
await config_backup_executor.execute(config, test_settings)
except Exception:
# May fail due to mocking complexity, but that's ok
pass
# Should have attempted to create tarfile
# assert mock_tar.called # Would check if it was actually called
@pytest.mark.asyncio
async def test_executor_handles_missing_source(self, test_settings: Settings, tmp_path: Path):
"""Test executor handles missing source directory."""
backup_dir = tmp_path / "backups"
backup_dir.mkdir()
config = {
"sources": [{
"path": "/nonexistent/path",
"name": "missing",
"excludes": []
}],
"backup_dir": str(backup_dir),
"compress": True
}
with patch('src.executors.config_backup_executor.Path') as mock_path_cls:
mock_source = MagicMock()
mock_source.exists.return_value = False
mock_path_cls.return_value = mock_source
result = await config_backup_executor.execute(config, test_settings)
# Should skip non-existent sources
assert "skipped" in result.lower() or "not found" in result.lower() or "0" in result
@pytest.mark.asyncio
async def test_cleanup_old_backups(self, tmp_path: Path):
"""Test cleanup of old backup files."""
backup_dir = tmp_path / "backups"
backup_dir.mkdir()
# Create some "old" backup files
old_backup = backup_dir / "backup-2020-01-01.tar.gz"
old_backup.write_text("old")
recent_backup = backup_dir / "backup-2025-12-01.tar.gz"
recent_backup.write_text("recent")
with patch('src.executors.config_backup_executor.Path') as mock_path_cls:
mock_backup_dir = MagicMock()
mock_old_file = MagicMock()
mock_old_file.name = "backup-2020-01-01.tar.gz"
mock_old_file.stat.return_value.st_mtime = 0 # Very old
mock_recent_file = MagicMock()
mock_recent_file.name = "backup-2025-12-01.tar.gz"
mock_recent_file.stat.return_value.st_mtime = 999999999999 # Recent
mock_backup_dir.glob.return_value = [mock_old_file, mock_recent_file]
mock_path_cls.return_value = mock_backup_dir
config_backup_executor._cleanup_old_backups(mock_backup_dir, retention_days=7)
# Old file should be removed
mock_old_file.unlink.assert_called_once()
assert not old.exists()
assert recent.exists()
assert unrelated.exists(), "cleanup must only touch files it wrote"
@pytest.mark.executor
@pytest.mark.unit
class TestConfigBackupHelpers:
"""Tests for helper functions."""
class TestEventLoopIsNotBlocked:
"""T-74. The scheduler serves its own API from the loop that runs executors.
def test_tar_filter_excludes_cache(self):
"""Test that tar filter excludes cache directories."""
excludes = ["*/cache/*", "*.log"]
filter_func = config_backup_executor._create_tar_filter(excludes)
This job spends ~21 minutes in tarfile and zlib, so calling it inline made
the whole service unreachable 03:05-03:25 every night. The hourly health
check runs at :35 and so never once observed it — the outage was invisible
for as long as it existed.
"""
# Mock tarinfo for cache file
cache_tarinfo = MagicMock()
cache_tarinfo.name = "data/cache/temp.txt"
@pytest.mark.asyncio
async def test_run_yields_to_the_loop_while_backing_up(
self, test_settings: Settings, monkeypatch
):
blocked_for = 0.4
monkeypatch.setattr(
config_backup_executor, "_backup",
lambda config: (time.sleep(blocked_for), "Backup completed: fake")[1],
)
result = filter_func(cache_tarinfo)
assert result is None # Should exclude
ticks = 0
def test_tar_filter_includes_normal_files(self):
"""Test that tar filter includes normal files."""
excludes = ["*/cache/*"]
filter_func = config_backup_executor._create_tar_filter(excludes)
async def heartbeat():
nonlocal ticks
while True:
await asyncio.sleep(0.01)
ticks += 1
# Mock tarinfo for normal file
normal_tarinfo = MagicMock()
normal_tarinfo.name = "data/config.json"
hb = asyncio.create_task(heartbeat())
try:
result = await config_backup_executor._run({}, test_settings)
finally:
hb.cancel()
result = filter_func(normal_tarinfo)
assert result == normal_tarinfo # Should include
assert result == "Backup completed: fake"
# Held inline, the loop gets no scheduling opportunity at all and this is
# 0. Off the loop it is ~40. The bar is low on purpose: the distinction
# being drawn is "the loop ran" versus "the loop was dead", and a loaded
# CI box should not turn that into a flake.
assert ticks >= 5, f"event loop starved during backup: {ticks} ticks"
def test_tar_filter_with_wildcard_patterns(self):
"""Test tar filter with various wildcard patterns."""
excludes = ["*.log", "*.tmp", "temp/*"]
filter_func = config_backup_executor._create_tar_filter(excludes)
# Log file
log_tarinfo = MagicMock()
log_tarinfo.name = "app.log"
assert filter_func(log_tarinfo) is None
@pytest.mark.executor
@pytest.mark.unit
class TestExcludeMatching:
"""Substring matching, and why it must stay that way.
# Temp file
tmp_tarinfo = MagicMock()
tmp_tarinfo.name = "cache.tmp"
assert filter_func(tmp_tarinfo) is None
`_create_tar_filter` reduces each pattern to a literal core and asks whether
it appears anywhere in the member name. That reads like an unfinished glob,
and the obvious "improvement" is `fnmatch`. These tests exist to make that
change fail loudly, because the deployed config is written against these
semantics and `fnmatch` would silently stop excluding the largest things in
the tree.
"""
# Normal file
normal_tarinfo = MagicMock()
normal_tarinfo.name = "config.json"
assert filter_func(normal_tarinfo) == normal_tarinfo
def test_mid_path_fragment_matches_anywhere(self):
"""`ollama/models/*` must exclude a member named `docker-data/ollama/...`.
Under fnmatch the pattern anchors at the start of the name, does not
match, and many GB of model blobs enter the nightly archive.
"""
f = config_backup_executor._create_tar_filter(["ollama/models/*"])
class TI:
name = "docker-data/ollama/models/blobs/sha256-abc"
assert f(TI()) is None
def test_bare_suffix_matches_every_file_carrying_it(self):
"""The deployment excludes logs by the bare string `.log`, not `*.log`.
Under fnmatch this matches only a file named exactly `.log`, so every
real log file starts being archived.
"""
f = config_backup_executor._create_tar_filter([".log"])
class TI:
name = "docker-data/svc/logs/app.log"
assert f(TI()) is None
def test_wrapped_pattern_is_reduced_to_its_core(self):
f = config_backup_executor._create_tar_filter(["*/cache/*"])
class Cache:
name = "docker-data/svc/cache/junk.bin"
class Normal:
name = "docker-data/svc/settings.json"
normal = Normal()
assert f(Cache()) is None
assert f(normal) is normal
def test_a_glob_star_is_not_interpreted(self):
"""`*.log` is a literal here — it is not a suffix match.
This is the sharp edge of substring matching and the reason the deployed
config spells the pattern `.log`. Pinned so the behaviour is documented
rather than discovered.
"""
f = config_backup_executor._create_tar_filter(["*.log"])
class TI:
name = "app.log"
ti = TI()
assert f(ti) is ti
def test_no_excludes_keeps_everything(self):
f = config_backup_executor._create_tar_filter([])
class TI:
name = "anything/at/all"
ti = TI()
assert f(ti) is ti
+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