16 Commits
Author SHA1 Message Date
jpmschweitzerandClaude 1c861e2fe1 release v1.6.0
Build and Push / release (push) Successful in 4s
Build and Push / build (push) Successful in 1m15s
Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-11 12:09:35 +02:00
jpmschweitzerandClaude e45de4fec7 fix(orchestrator): record the terminal states that were never written
Two defects with one shape: an execution reaches a terminal condition and the
orchestrator fails to write it down, so the system's own record disagrees with
what happened. Neither produced an error. Both produced silence.

T-2 -- a restart mid-task unscheduled that task forever.

get_tasks_for_minute excludes any task holding a task_executions row with
status='running'. The row is written before the executor runs and updated
after, so a process dying in between left it 'running' permanently, and the
task was then excluded from every future minute with no error, no alarm and no
log line. It did not fail; it went quiet.

test_example_task had held such a row since 2025-12-07 -- 5916 hours. It is
disabled, so nothing was broken by that instance; the mechanism is the point,
and the exposure is daily, because Watchtower restarts this container at 4 AM
while the config backup starts at 03:05 and runs ~21 minutes.

Startup now reconciles them, where the reasoning is sound by construction: this
process has just begun, so nothing it can see is genuinely running.

Marked 'orphaned', not 'failed'. When the process dies mid-task the work may
well have completed -- a backup that finished and never got to update its row
is indistinguishable from one that died halfway -- and 'failed' would assert an
outcome nobody observed. Same error as the health-report diagnostic fixed in
18be804: naming a cause you did not witness.

Not extended to a duration-based sweep. While this process lives, execute_task's
finally clause always closes the row, so a stale row implies a dead owner. A
time-based rule would have to tell a slow task from a dead one, and getting that
wrong closes the record of a task still working.

T-3 -- the 'timeout' status was unreachable.

execute_task has an `except asyncio.TimeoutError` branch that records
status='timeout'. It could never run: _run_executor wrapped the awaited call in
`except Exception`, and since 3.11 asyncio.TimeoutError IS the builtin
TimeoutError (OSError -> Exception), so the broad handler caught it first and
converted it to an ordinary error tuple. Confirmed in the deployed runtime and
against the history -- 18,785 executions since 2025-12-07, of which 'timeout'
rows: zero. Every timeout in eight months was filed as a generic failure,
erasing the distinction between "too slow for its window" and "broken".

A narrower except after a broader one is dead code, and no linter is configured
here to say so.

One trap in fixing it: while the branch was unreachable a timeout travelled the
normal path, which DOES update scheduled_tasks. Making the branch reachable
without that write would have traded a wrong status for a stale one, so
_update_task_outcome now mirrors terminal outcomes onto the parent row.

The timeout message also states that the work may still be running -- after
T-74 executors are handed to asyncio.to_thread, and a thread cannot be
cancelled, so wait_for frees the loop while the work continues.

Both fixes are mutation-checked: removing the startup call fails the ordering
test, removing the narrow except clause fails the propagation test. Suite goes
118 -> 126 passing with the same 36 pre-existing failures.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-11 12:09:35 +02:00
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
27 changed files with 2065 additions and 312 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
+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);
+47
View File
@@ -0,0 +1,47 @@
INSERT INTO ticket_history (ticket_record_id, field, old_value, new_value, changed_by, changed_at, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0FHPBBXXKBRWGR3FJZRK58', 'description', NULL, 'DELETE /tasks/{task_name} returns 500 whenever the task has at least one row in task_executions:
psycopg2.errors.ForeignKeyViolation: update or delete on table "scheduled_tasks"
violates foreign key constraint "task_executions_task_id_fkey" on table "task_executions"
DETAIL: Key (id)=(46) is still referenced from table "task_executions".
src/main.py:355. Since every task that has ever fired has execution history, the endpoint works only for tasks that have never run — which is close to none of them. Found on 2026-08-11 while cleaning up a temporary probe task created for T-74; it had to be removed with hand-written SQL against two tables, which is not something the API should require.
The caller gets a bare "Internal Server Error" with no indication that history is the obstacle, so it reads as the service being broken rather than the request being refusable.
Deciding what delete should MEAN is the actual work here, and it should not be guessed:
- cascade — drop the execution history with the task. Simple, and silently destroys the audit trail for a task someone deletes by mistake.
- soft delete — mark it deleted and keep the history. Keeps the audit trail, adds a state every query then has to filter on.
- refuse with 409 and a real message ("task has N executions; pass ?purge=true"). Explicit, and makes the destructive variant a deliberate act.
The third is the smallest correct change and matches how the rest of this system treats destructive operations. Whichever is chosen, a 500 on a foreseeable, well-defined condition is the part that is simply wrong.', NULL, '2026-08-11 09:53:56', '2026-08-11 09:53:56.969', '2026-08-11 09:53:56.969', NULL, '8b2dfce5342a7cd89bbe6d04c03157f1', 2) ON CONFLICT(hash) DO NOTHING;
INSERT INTO ticket_history (ticket_record_id, field, old_value, new_value, changed_by, changed_at, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0GQAR5KXM3PTG0QT3R1PV8', 'description', NULL, 'get_tasks_for_minute (src/tasks/executor.py:70) excludes any task holding a task_executions row with status=''running'':
AND id NOT IN (SELECT task_id FROM task_executions WHERE status = ''running'')
Nothing ever reconciles that row. It is written before the executor runs and updated after, so a process that dies in between leaves it ''running'' permanently — and the task is then excluded from every future minute, forever, with no error, no alarm and no log line. The task simply stops existing as far as the scheduler is concerned.
This is not hypothetical. test_example_task has held a ''running'' row since 2025-12-07 — 5916 hours. It happens to be disabled, so nothing is broken today; the mechanism is what matters, not this instance.
The exposure is real and daily: Watchtower restarts this container every morning at 4 AM. Any task still running at that moment is permanently unscheduled by it. The config backup runs 03:05 and takes ~21 minutes, which is not far off.
The failure mode is the one this system keeps producing: it looks like nothing. A backup that stops running forever produces no failure — it produces silence, and silence reads as health.
FIX: reconcile at startup. Any row left ''running'' when the process starts cannot be running, because the process that owned it is gone. Mark those ''orphaned'' (or ''failed'' with a message naming the restart) during app startup, before the scheduler begins its first minute. Consider also a stale-row guard for rows older than the task''s timeout_seconds, which covers a killed worker without a restart.
The audit trail is preserved either way — this is a status correction, not a deletion.', NULL, '2026-08-11 09:59:05', '2026-08-11 09:59:05.288', '2026-08-11 09:59:05.288', NULL, '1b14bc4af4acc6df0dc2d21675edf31f', 2) ON CONFLICT(hash) DO NOTHING;
INSERT INTO ticket_history (ticket_record_id, field, old_value, new_value, changed_by, changed_at, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0GQBTCRCTXK6EJF0T8CS2C', 'description', NULL, 'execute_task has an `except asyncio.TimeoutError` handler (src/tasks/executor.py:179) that records status=''timeout''. It can never run.
_run_executor wraps the awaited call in `except Exception` (line 219) and returns the error as a value. Since Python 3.11 asyncio.TimeoutError IS the builtin TimeoutError, which inherits OSError -> Exception, so the broad handler catches it first and converts it into an ordinary (None, error) tuple. execute_task then sees a non-None error and files status=''failed''.
Verified in the deployed runtime:
asyncio.TimeoutError is TimeoutError: True
MRO: TimeoutError -> OSError -> Exception -> BaseException
And confirmed against 8 months of history — 18,785 executions, and the count of status=''timeout'' rows is zero:
success 16690 | failed 2093 | running 2 | timeout 0
So every timeout since 2025-12-07 has been recorded as a generic failure carrying a traceback. Operationally that erases the distinction that matters most when a job misbehaves: "this job is too slow for its window" and "this job is broken" need different responses, and right now they look identical in the history.
FIX: catch asyncio.TimeoutError explicitly in _run_executor, ahead of the broad handler, and let it propagate (or return a marker execute_task can distinguish). Note the ordering is the whole bug — a narrower except after a broader one is dead code, and there is no linter configured here to say so.
Related: a thread started by asyncio.to_thread cannot be cancelled, so after T-74 a timeout frees the loop while the work continues to completion. Whatever ''timeout'' comes to mean should say so rather than implying the work stopped.', NULL, '2026-08-11 09:59:05', '2026-08-11 09:59:05.545', '2026-08-11 09:59:05.545', NULL, 'af08842306d6326036b1865b8cbf7c58', 2) ON CONFLICT(hash) DO NOTHING;
+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);
+3
View File
@@ -0,0 +1,3 @@
INSERT INTO ticket_idmap (record_id, ticket_id, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0FHPBBXXKBRWGR3FJZRK58', 'T-1', '2026-08-11 09:53:56.835', '2026-08-11 09:53:56.835', NULL, '3d980945e785a6bc7ca8fcaa8250e22b', 2) ON CONFLICT(record_id) DO UPDATE SET ticket_id=excluded.ticket_id, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at >= ticket_idmap.updated_at;
INSERT INTO ticket_idmap (record_id, ticket_id, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0GQAR5KXM3PTG0QT3R1PV8', 'T-2', '2026-08-11 09:59:05.153', '2026-08-11 09:59:05.153', NULL, 'c8655ec9601ba93fe822395329e62262', 2) ON CONFLICT(record_id) DO UPDATE SET ticket_id=excluded.ticket_id, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at >= ticket_idmap.updated_at;
INSERT INTO ticket_idmap (record_id, ticket_id, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0GQBTCRCTXK6EJF0T8CS2C', 'T-3', '2026-08-11 09:59:05.427', '2026-08-11 09:59:05.427', NULL, '3bce28c3d803d3f5036cfbb1ac11969c', 2) ON CONFLICT(record_id) DO UPDATE SET ticket_id=excluded.ticket_id, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at >= ticket_idmap.updated_at;
@@ -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);
+50
View File
@@ -0,0 +1,50 @@
INSERT INTO tickets (record_id, type, parent_record_id, title, description, status, priority, assigned_to, team, decision_ref, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0FHPBBXXKBRWGR3FJZRK58', 'bug', NULL, 'DELETE /tasks/{name} 500s for any task that has ever run', NULL, 'backlog', 'high', NULL, NULL, NULL, '2026-08-11 09:53:56.826', '2026-08-11 09:53:56.826', NULL, '86310746d98f53b722c7336b0bab980a', 2) ON CONFLICT(record_id) DO UPDATE SET type=excluded.type, parent_record_id=excluded.parent_record_id, title=excluded.title, description=excluded.description, status=excluded.status, priority=excluded.priority, assigned_to=excluded.assigned_to, team=excluded.team, decision_ref=excluded.decision_ref, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at >= tickets.updated_at;
INSERT INTO tickets (record_id, type, parent_record_id, title, description, status, priority, assigned_to, team, decision_ref, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0FHPBBXXKBRWGR3FJZRK58', 'bug', NULL, 'DELETE /tasks/{name} 500s for any task that has ever run', 'DELETE /tasks/{task_name} returns 500 whenever the task has at least one row in task_executions:
psycopg2.errors.ForeignKeyViolation: update or delete on table "scheduled_tasks"
violates foreign key constraint "task_executions_task_id_fkey" on table "task_executions"
DETAIL: Key (id)=(46) is still referenced from table "task_executions".
src/main.py:355. Since every task that has ever fired has execution history, the endpoint works only for tasks that have never run — which is close to none of them. Found on 2026-08-11 while cleaning up a temporary probe task created for T-74; it had to be removed with hand-written SQL against two tables, which is not something the API should require.
The caller gets a bare "Internal Server Error" with no indication that history is the obstacle, so it reads as the service being broken rather than the request being refusable.
Deciding what delete should MEAN is the actual work here, and it should not be guessed:
- cascade — drop the execution history with the task. Simple, and silently destroys the audit trail for a task someone deletes by mistake.
- soft delete — mark it deleted and keep the history. Keeps the audit trail, adds a state every query then has to filter on.
- refuse with 409 and a real message ("task has N executions; pass ?purge=true"). Explicit, and makes the destructive variant a deliberate act.
The third is the smallest correct change and matches how the rest of this system treats destructive operations. Whichever is chosen, a 500 on a foreseeable, well-defined condition is the part that is simply wrong.', 'backlog', 'high', NULL, NULL, NULL, '2026-08-11 09:53:56.826', '2026-08-11 09:53:56.969', NULL, '248bb53ca2bae0f12fabd9d88a1ca171', 2) ON CONFLICT(record_id) DO UPDATE SET type=excluded.type, parent_record_id=excluded.parent_record_id, title=excluded.title, description=excluded.description, status=excluded.status, priority=excluded.priority, assigned_to=excluded.assigned_to, team=excluded.team, decision_ref=excluded.decision_ref, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at >= tickets.updated_at;
INSERT INTO tickets (record_id, type, parent_record_id, title, description, status, priority, assigned_to, team, decision_ref, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0GQAR5KXM3PTG0QT3R1PV8', 'bug', NULL, 'A restart mid-task unschedules that task forever, silently', NULL, 'backlog', 'high', NULL, NULL, NULL, '2026-08-11 09:59:05.153', '2026-08-11 09:59:05.153', NULL, '0bda80455ba51ea4cb0cb91e302ac216', 2) ON CONFLICT(record_id) DO UPDATE SET type=excluded.type, parent_record_id=excluded.parent_record_id, title=excluded.title, description=excluded.description, status=excluded.status, priority=excluded.priority, assigned_to=excluded.assigned_to, team=excluded.team, decision_ref=excluded.decision_ref, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at >= tickets.updated_at;
INSERT INTO tickets (record_id, type, parent_record_id, title, description, status, priority, assigned_to, team, decision_ref, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0GQAR5KXM3PTG0QT3R1PV8', 'bug', NULL, 'A restart mid-task unschedules that task forever, silently', 'get_tasks_for_minute (src/tasks/executor.py:70) excludes any task holding a task_executions row with status=''running'':
AND id NOT IN (SELECT task_id FROM task_executions WHERE status = ''running'')
Nothing ever reconciles that row. It is written before the executor runs and updated after, so a process that dies in between leaves it ''running'' permanently — and the task is then excluded from every future minute, forever, with no error, no alarm and no log line. The task simply stops existing as far as the scheduler is concerned.
This is not hypothetical. test_example_task has held a ''running'' row since 2025-12-07 — 5916 hours. It happens to be disabled, so nothing is broken today; the mechanism is what matters, not this instance.
The exposure is real and daily: Watchtower restarts this container every morning at 4 AM. Any task still running at that moment is permanently unscheduled by it. The config backup runs 03:05 and takes ~21 minutes, which is not far off.
The failure mode is the one this system keeps producing: it looks like nothing. A backup that stops running forever produces no failure — it produces silence, and silence reads as health.
FIX: reconcile at startup. Any row left ''running'' when the process starts cannot be running, because the process that owned it is gone. Mark those ''orphaned'' (or ''failed'' with a message naming the restart) during app startup, before the scheduler begins its first minute. Consider also a stale-row guard for rows older than the task''s timeout_seconds, which covers a killed worker without a restart.
The audit trail is preserved either way — this is a status correction, not a deletion.', 'backlog', 'high', NULL, NULL, NULL, '2026-08-11 09:59:05.153', '2026-08-11 09:59:05.288', NULL, '23a72bff9ffc3592c796741df8a7232e', 2) ON CONFLICT(record_id) DO UPDATE SET type=excluded.type, parent_record_id=excluded.parent_record_id, title=excluded.title, description=excluded.description, status=excluded.status, priority=excluded.priority, assigned_to=excluded.assigned_to, team=excluded.team, decision_ref=excluded.decision_ref, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at >= tickets.updated_at;
INSERT INTO tickets (record_id, type, parent_record_id, title, description, status, priority, assigned_to, team, decision_ref, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0GQBTCRCTXK6EJF0T8CS2C', 'bug', NULL, 'The ''timeout'' execution status is unreachable; timeouts are filed as generic failures', NULL, 'backlog', 'medium', NULL, NULL, NULL, '2026-08-11 09:59:05.427', '2026-08-11 09:59:05.427', NULL, '31a32d6dce87263657892d48552d1f0a', 2) ON CONFLICT(record_id) DO UPDATE SET type=excluded.type, parent_record_id=excluded.parent_record_id, title=excluded.title, description=excluded.description, status=excluded.status, priority=excluded.priority, assigned_to=excluded.assigned_to, team=excluded.team, decision_ref=excluded.decision_ref, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at >= tickets.updated_at;
INSERT INTO tickets (record_id, type, parent_record_id, title, description, status, priority, assigned_to, team, decision_ref, created_at, updated_at, deleted_at, hash, canonical_version) VALUES ('06FZ0GQBTCRCTXK6EJF0T8CS2C', 'bug', NULL, 'The ''timeout'' execution status is unreachable; timeouts are filed as generic failures', 'execute_task has an `except asyncio.TimeoutError` handler (src/tasks/executor.py:179) that records status=''timeout''. It can never run.
_run_executor wraps the awaited call in `except Exception` (line 219) and returns the error as a value. Since Python 3.11 asyncio.TimeoutError IS the builtin TimeoutError, which inherits OSError -> Exception, so the broad handler catches it first and converts it into an ordinary (None, error) tuple. execute_task then sees a non-None error and files status=''failed''.
Verified in the deployed runtime:
asyncio.TimeoutError is TimeoutError: True
MRO: TimeoutError -> OSError -> Exception -> BaseException
And confirmed against 8 months of history — 18,785 executions, and the count of status=''timeout'' rows is zero:
success 16690 | failed 2093 | running 2 | timeout 0
So every timeout since 2025-12-07 has been recorded as a generic failure carrying a traceback. Operationally that erases the distinction that matters most when a job misbehaves: "this job is too slow for its window" and "this job is broken" need different responses, and right now they look identical in the history.
FIX: catch asyncio.TimeoutError explicitly in _run_executor, ahead of the broad handler, and let it propagate (or return a marker execute_task can distinguish). Note the ordering is the whole bug — a narrower except after a broader one is dead code, and there is no linter configured here to say so.
Related: a thread started by asyncio.to_thread cannot be cancelled, so after T-74 a timeout frees the loop while the work continues to completion. Whatever ''timeout'' comes to mean should say so rather than implying the work stopped.', 'backlog', 'medium', NULL, NULL, NULL, '2026-08-11 09:59:05.427', '2026-08-11 09:59:05.545', NULL, 'ae0db866086e38b681b0ea32837df277', 2) ON CONFLICT(record_id) DO UPDATE SET type=excluded.type, parent_record_id=excluded.parent_record_id, title=excluded.title, description=excluded.description, status=excluded.status, priority=excluded.priority, assigned_to=excluded.assigned_to, team=excluded.team, decision_ref=excluded.decision_ref, updated_at=excluded.updated_at, deleted_at=excluded.deleted_at, hash=excluded.hash, canonical_version=excluded.canonical_version WHERE excluded.updated_at >= tickets.updated_at;
-72
View File
@@ -1,72 +0,0 @@
# AGENTS.md
> **Start every session by reading this file.**
> This file outlines the operational protocols, coding standards, and architectural decisions for this FastAPI project.
## 1. Agent Operational Protocols
### 🧠 Work Patterns (Plan-Act-Reflect)
* **Plan:** Before writing code, briefly outline your plan. Identify which files you will touch and what the side effects might be.
* **Act:** Execute the changes in small, atomic steps.
* **Reflect:** After coding, verify your work. Did you break existing tests? Did you add new tests?
### 🛡️ Git Discipline
* **NEVER commit to `main` or `master` directly.** Always create a feature branch: `feature/your-feature-name` or `fix/issue-description`.
* **Commit Messages:** Use the [Conventional Commits](https://www.conventionalcommits.org/) format.
* `feat: add user login endpoint`
* `fix: resolve database connection timeout`
* `refactor: split monolith dependency file`
* **Atomic Commits:** Keep commits small. One logical change = one commit.
### 📝 Changelog Maintenance
* **Update `CHANGELOG.md`** with every user-facing change.
* Format: `## [Unreleased] - YYYY-MM-DD` followed by `### Added`, `### Changed`, or `### Fixed`.
### 🚀 Release Flow
When changes are ready for deployment:
1. **Ask user if deploy cycle is desired **
2. **Update version** in `pyproject.toml`:
- Bug fixes: bump patch version (1.8.3 → 1.8.4)
- New features: bump minor version (1.8.4 → 1.9.0)
3. **Update CHANGELOG.md**:
- Move items from `[Unreleased]` to new version section
- Add release date: `## [1.8.4] - 2025-12-16`
4. **Commit and tag**:
```bash
git add -A
git commit -m "fix: description of changes"
git tag v1.8.4
git push origin main --tags
```
5. **CI/CD triggers automatically**:
- Gitea CI builds Docker image on new tag
- Watchtower pulls and deploys to production
- Verify deployment: `curl http://192.168.86.149:8000/health`
---
## 2. FastAPI Architecture & Best Practices
*Reference: [FastAPI Best Practices](https://github.com/zhanymkanov/fastapi-best-practices)*
### 📂 Project Structure (Directory-based, NOT File-type based)
Do **not** group files by type (e.g., one huge `routers` folder). Group by **domain/module** inside a `src/` directory.
**Correct Structure:**
```text
src/
├── auth/
│ ├── router.py # Endpoints
│ ├── schemas.py # Pydantic models
│ ├── service.py # Business logic (CRUD, etc.)
│ ├── dependencies.py# Module-specific dependencies
│ └── config.py # Module-specific settings
├── posts/
│ ├── router.py
│ └── ...
└── main.py # App entry point
+38 -3
View File
@@ -6,6 +6,43 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/).
## [Unreleased]
## [1.6.0] - 2026-08-11
### Fixed
- A restart mid-task no longer unschedules that task forever. An execution row left
`running` excluded its task from every future minute, silently; startup now releases them.
- Timeouts are recorded as `timeout` instead of a generic failure. The status existed but was
unreachable, so all 18,785 executions since December contain zero of them.
### Added
- `orphaned` execution status — an execution whose process died, whose outcome is unknown.
Distinct from `failed`, which asserts the work did not succeed.
## [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
@@ -269,9 +306,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.6.0` (verified 2026-08-11). 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)"
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)_
+1 -1
View File
@@ -1,6 +1,6 @@
[project]
name = "the-scheduler"
version = "1.4.0"
version = "1.6.0"
description = "System-wide maintenance orchestration - backups, doc mirroring, cleanup, task automation"
readme = "README.md"
requires-python = ">=3.12"
+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
+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")
+32 -1
View File
@@ -37,6 +37,7 @@ from pathlib import Path
import httpx
from src.config import Settings
from src.executors import health_report
logger = logging.getLogger(__name__)
@@ -72,7 +73,7 @@ def _prune(output_dir: Path, retention_days: int) -> int:
return removed
async def execute(config: dict, settings: Settings) -> str:
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"))
@@ -138,3 +139,33 @@ async def execute(config: dict, settings: Settings) -> str:
)
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
+7
View File
@@ -65,6 +65,13 @@ async def lifespan(app: FastAPI):
task_executor = TaskExecutor(settings)
logger.info("Task executor initialized (max 5 concurrent tasks)")
# Before the first minute is processed, release any task still held by an
# execution row belonging to an instance that no longer exists. A 'running'
# row excludes its task from scheduling permanently, so skipping this leaves
# tasks silently unschedulable across every restart. Blocking briefly is fine
# here — the app serves no requests until lifespan yields.
task_executor.reconcile_orphaned_executions()
# Initialize APScheduler with minimal configuration
# No jobstore needed - we only have one in-memory job
scheduler = AsyncIOScheduler(
+116 -2
View File
@@ -115,6 +115,70 @@ class TaskExecutor:
return True
def reconcile_orphaned_executions(self) -> int:
"""Close out execution rows left 'running' by a process that is gone.
get_tasks_for_minute excludes any task holding a 'running' row. That row
is written before the executor runs and updated after, so a process that
dies in between leaves it 'running' forever — and the task is then
excluded from every future minute, permanently, with no error and no log
line. It does not fail; it goes quiet, and quiet reads as healthy.
Live exposure rather than theory: Watchtower restarts this container at
4 AM daily, and the config backup starts at 03:05 and runs ~21 minutes.
A row from 2025-12-07 sat 'running' for eight months before anyone
looked.
Called at startup, where the reasoning is sound by construction: this
process has just begun, so nothing it can see is genuinely running, and
any such row belongs to an instance that no longer exists.
Marked 'orphaned', not 'failed'. When the process dies mid-task the work
may well have finished — a backup that completed and never got to update
its row is indistinguishable from one that died halfway. 'failed' would
assert an outcome nobody observed. 'orphaned' says only what is known:
we lost track of it.
Deliberately not extended to a time-based sweep of long-running rows.
While this process lives, execute_task's finally clause always closes the
row out, so a stale row implies a dead owner. A duration-based rule would
have to tell a slow task from a dead one, and getting that wrong closes
the record of a task that is still working.
"""
try:
with self.get_db_connection() as conn:
with conn.cursor() as cur:
cur.execute("""
UPDATE task_executions
SET status = 'orphaned',
completed_at = %s,
error = 'Scheduler restarted while this execution was '
'running; its outcome is unknown.'
WHERE status = 'running'
RETURNING task_name, started_at
""", (datetime.now(timezone.utc),))
orphans = cur.fetchall()
conn.commit()
except Exception as e:
# Never fatal. A scheduler that refuses to start because it could not
# tidy up is worse than one carrying a stale row.
logger.error(f"Could not reconcile orphaned executions: {e}")
return 0
for task_name, started_at in orphans:
logger.warning(
f"Orphaned execution recovered: {task_name} was left 'running' "
f"since {started_at}. That task had been excluded from scheduling "
f"until now."
)
if orphans:
logger.warning(
f"{len(orphans)} task(s) were unschedulable and are now released."
)
else:
logger.info("No orphaned executions to reconcile")
return len(orphans)
async def execute_task(self, task: Dict[str, Any]):
"""
Execute a single task with timeout and error handling.
@@ -178,8 +242,17 @@ class TaskExecutor:
except asyncio.TimeoutError:
logger.error(f"Task {task_name} timed out after {timeout}s")
self._update_execution_status(execution_id, 'timeout',
error=f"Task exceeded timeout of {timeout}s")
self._update_execution_status(
execution_id, 'timeout',
error=f"Task exceeded timeout of {timeout}s. The underlying work may "
f"still be running — executors that use asyncio.to_thread hand "
f"the work to a thread, and a thread cannot be cancelled.")
# scheduled_tasks has to be written here as well. On the normal path
# it is updated alongside the execution row, and while this branch was
# unreachable a timeout travelled that path as a 'failed' result — so
# last_status did stay current. Making the branch reachable without
# this call would swap one wrong status for a stale one.
self._update_task_outcome(task_id, 'timeout', started_at)
except Exception as e:
logger.error(f"Task {task_name} failed with exception: {e}")
@@ -214,11 +287,52 @@ class TaskExecutor:
return result, None
except asyncio.TimeoutError:
# This clause must precede `except Exception`, and that ordering is
# the entire bug it fixes. Since Python 3.11 asyncio.TimeoutError IS
# the builtin TimeoutError, which inherits OSError -> Exception, so
# the broad handler below used to catch it first and convert it into
# an ordinary (None, error) tuple. execute_task then filed it as a
# generic 'failed', and its own `except asyncio.TimeoutError` branch
# was unreachable: zero 'timeout' rows across 18,785 executions and
# eight months of history.
#
# Re-raised rather than returned, because the distinction is the
# point: "too slow for its window" and "broken" call for different
# responses and were indistinguishable in the record.
raise
except ModuleNotFoundError:
return None, f"Executor module not found: {executor_name}"
except Exception as e:
return None, f"Executor error: {str(e)}\n{traceback.format_exc()}"
def _update_task_outcome(self, task_id: int, status: str, started_at: datetime):
"""Mirror a terminal outcome onto scheduled_tasks.
The happy path writes task_executions and scheduled_tasks in one
transaction. The error branches historically wrote only the former, which
did not show while every timeout was being funnelled through the happy
path as a 'failed'. Once a branch bypasses that path it has to keep
last_run/last_status current itself, or the task list quietly reports the
previous run's outcome as though it were the latest.
"""
try:
completed_at = datetime.now(timezone.utc)
with self.get_db_connection() as conn:
with conn.cursor() as cur:
cur.execute("""
UPDATE scheduled_tasks
SET last_run = %s, last_status = %s,
last_duration_seconds = %s, updated_at = %s
WHERE id = %s
""", (completed_at, status,
int((completed_at - started_at).total_seconds()),
completed_at, task_id))
conn.commit()
except Exception as e:
logger.error(f"Failed to update task outcome for {task_id}: {e}")
def _update_execution_status(self, execution_id: int, status: str, error: str = None):
"""Update execution record with final status."""
if execution_id is None:
+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
+178
View File
@@ -305,3 +305,181 @@ class TestTaskExecutor:
# Should only execute max_concurrent (5) tasks
assert mock_execute.call_count <= executor.max_concurrent
@pytest.mark.unit
class TestOrphanReconciliation:
"""T-2. A 'running' row excludes its task from scheduling forever.
get_tasks_for_minute filters out any task holding one, and nothing ever
closed those rows, so a process that died between writing the row and
updating it left its task permanently unschedulable — silently. A row from
2025-12-07 sat that way for eight months. Watchtower restarts this container
nightly, so the exposure was daily.
"""
def _mock_conn(self, executor, rows):
conn, cur = MagicMock(), MagicMock()
cur.fetchall.return_value = rows
conn.cursor.return_value.__enter__ = MagicMock(return_value=cur)
conn.cursor.return_value.__exit__ = MagicMock(return_value=None)
conn.__enter__ = MagicMock(return_value=conn)
conn.__exit__ = MagicMock(return_value=None)
patcher = patch.object(executor, 'get_db_connection', return_value=conn)
patcher.start()
return cur, patcher
def test_running_rows_are_released(self, test_settings: Settings):
executor = TaskExecutor(test_settings)
cur, p = self._mock_conn(executor, [("backup_docker_configs_daily", datetime(2026, 8, 11))])
try:
assert executor.reconcile_orphaned_executions() == 1
sql = cur.execute.call_args[0][0]
assert "UPDATE task_executions" in sql
# It must target exactly the rows get_tasks_for_minute excludes, and
# move them to a status it does not exclude. If these two ever drift
# apart the bug returns in silence.
assert "WHERE status = 'running'" in sql
assert "status = 'orphaned'" in sql
finally:
p.stop()
def test_clean_startup_reports_nothing(self, test_settings: Settings):
executor = TaskExecutor(test_settings)
cur, p = self._mock_conn(executor, [])
try:
assert executor.reconcile_orphaned_executions() == 0
finally:
p.stop()
def test_a_database_failure_does_not_stop_startup(self, test_settings: Settings):
"""Refusing to boot because cleanup failed is worse than a stale row."""
executor = TaskExecutor(test_settings)
with patch.object(executor, 'get_db_connection', side_effect=Exception("db down")):
assert executor.reconcile_orphaned_executions() == 0 # no raise
def test_orphaned_is_not_failed(self, test_settings: Settings):
"""The outcome is unknown, not known-bad.
A backup that finished and never got to update its row looks identical to
one that died halfway. Recording 'failed' asserts something nobody
observed.
"""
executor = TaskExecutor(test_settings)
cur, p = self._mock_conn(executor, [("t", datetime(2026, 1, 1))])
try:
executor.reconcile_orphaned_executions()
sql = cur.execute.call_args[0][0]
assert "'failed'" not in sql
finally:
p.stop()
@pytest.mark.asyncio
async def test_startup_reconciles_before_the_scheduler_starts(
self, test_settings: Settings, monkeypatch
):
"""Order matters: reconcile must finish before the first minute is processed.
Run the other way round and the first tick still sees the stale rows.
"""
from src import main
order = []
monkeypatch.setattr(main, 'get_settings', lambda: test_settings)
monkeypatch.setattr(
TaskExecutor, 'reconcile_orphaned_executions',
lambda self: (order.append('reconcile'), 0)[1],
)
class FakeScheduler:
def add_job(self, **kw): order.append('add_job')
def start(self): order.append('start')
def shutdown(self, wait=True): order.append('shutdown')
monkeypatch.setattr(main, 'AsyncIOScheduler', lambda **kw: FakeScheduler())
async with main.lifespan(None):
pass
assert 'reconcile' in order, "startup never reconciled orphaned executions"
assert order.index('reconcile') < order.index('start')
@pytest.mark.unit
class TestTimeoutIsDistinguishable:
"""T-3. The 'timeout' status existed in the code and had never been written.
_run_executor's `except Exception` sat above execute_task's
`except asyncio.TimeoutError`, and since 3.11 asyncio.TimeoutError IS the
builtin TimeoutError (OSError -> Exception), so the broad handler always won.
Eight months, 18,785 executions, zero timeout rows — every one filed as a
generic failure, erasing the difference between "too slow for its window" and
"broken".
"""
@pytest.mark.asyncio
async def test_timeout_propagates_instead_of_becoming_an_error_tuple(
self, test_settings: Settings, monkeypatch
):
import sys, types, asyncio as aio
mod = types.ModuleType("src.executors.slow_probe")
async def execute(config, settings):
await aio.sleep(5)
mod.execute = execute
monkeypatch.setitem(sys.modules, "src.executors.slow_probe", mod)
executor = TaskExecutor(test_settings)
with pytest.raises(aio.TimeoutError):
await executor._run_executor("slow_probe", {"config": {}}, timeout=0.05)
@pytest.mark.asyncio
async def test_a_timed_out_task_is_recorded_as_timeout(self, test_settings: Settings):
import asyncio as aio
executor = TaskExecutor(test_settings)
task = {'id': 7, 'task_name': 'slow', 'executor': 'slow_probe',
'service': 'scheduler', 'priority': 5, 'timeout_seconds': 1}
conn, cur = MagicMock(), MagicMock()
cur.fetchone.return_value = [123]
conn.cursor.return_value.__enter__ = MagicMock(return_value=cur)
conn.cursor.return_value.__exit__ = MagicMock(return_value=None)
conn.__enter__ = MagicMock(return_value=conn)
conn.__exit__ = MagicMock(return_value=None)
with patch.object(executor, 'get_db_connection', return_value=conn), \
patch.object(executor, '_run_executor',
new=AsyncMock(side_effect=aio.TimeoutError())), \
patch.object(executor, '_update_execution_status') as upd_exec, \
patch.object(executor, '_update_task_outcome') as upd_task:
await executor.execute_task(task)
assert upd_exec.call_args[0][1] == 'timeout', "execution row must say timeout"
# The trap: while the timeout branch was unreachable a timeout travelled
# the normal path, which DOES update scheduled_tasks. Making the branch
# reachable without this call would swap a wrong status for a stale one.
assert upd_task.called, "scheduled_tasks left stale after a timeout"
assert upd_task.call_args[0][1] == 'timeout'
@pytest.mark.asyncio
async def test_ordinary_errors_are_still_returned_not_raised(
self, test_settings: Settings, monkeypatch
):
"""The narrow clause must not swallow anything else on its way past."""
import sys, types
mod = types.ModuleType("src.executors.boom_probe")
async def execute(config, settings):
raise ValueError("kaboom")
mod.execute = execute
monkeypatch.setitem(sys.modules, "src.executors.boom_probe", mod)
executor = TaskExecutor(test_settings)
output, error = await executor._run_executor("boom_probe", {"config": {}}, timeout=5)
assert output is None
assert "kaboom" in error