Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1c861e2fe1 | ||
|
|
e45de4fec7 | ||
|
|
d91d0f2c63 | ||
|
|
b62024b21d | ||
|
|
18be804ba0 | ||
|
|
1012f6f374 | ||
|
|
d9cbaee1fc | ||
|
|
0c2199667f | ||
|
|
054974c3fb | ||
|
|
b712caefb1 | ||
|
|
5b4814a3cd | ||
|
|
9eb9e4b32b | ||
|
|
c64412d9b3 | ||
|
|
4cbfe7a6f9 | ||
|
|
e975f5b720 | ||
|
|
55fe0f9083 |
@@ -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:*)"
|
||||
]
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1 @@
|
||||
.pql/changelog/*.sql merge=union
|
||||
Executable
+13
@@ -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
@@ -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
|
||||
|
||||
@@ -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
|
||||
@@ -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);
|
||||
@@ -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;
|
||||
@@ -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,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);
|
||||
@@ -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,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;
|
||||
@@ -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
@@ -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)
|
||||
|
||||
@@ -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`.
|
||||
@@ -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
@@ -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"
|
||||
@@ -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
@@ -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"
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
@@ -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
|
||||
|
||||
@@ -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
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user