refactor(kanban_db_*): hand-compact docstrings/comments keeping every invariant; heartbeat SQL built once (4476->4390 LOC)
This commit is contained in:
@@ -160,20 +160,15 @@ def _cross_process_init_lock(path: Path):
|
||||
@contextlib.contextmanager
|
||||
def _dispatch_tick_lock(db_path: Path):
|
||||
"""Non-blocking single-writer guard around one dispatcher tick; yields
|
||||
``True`` if this process holds the board's dispatch lock, else ``False``
|
||||
(caller skips the tick).
|
||||
``True`` if this process holds the board's ``.dispatch.lock``, else
|
||||
``False`` (caller skips the tick).
|
||||
|
||||
An orphan gateway (``gateway run --replace`` / restart escaping a systemd
|
||||
or launchd service cgroup) becomes a second long-lived writer; two
|
||||
dispatchers both pass ``busy_timeout`` and race on WAL frames — the root
|
||||
cause of multi-writer corruption. ``_guard_supervised_gateway_conflict``
|
||||
blocks the common birth of an orphan; this lock is defense-in-depth.
|
||||
|
||||
**Non-blocking** on purpose: the gateway's async watcher must never stall;
|
||||
a loser retries next interval while the winner makes progress. Board-scoped
|
||||
via a ``.dispatch.lock`` sibling of ``kanban.db``. Without
|
||||
``fcntl``/``msvcrt`` it degrades to a no-op (yields ``True``) — the orphan
|
||||
scenario is specific to POSIX service managers.
|
||||
Two dispatchers (e.g. an orphan gateway escaping its service cgroup) both
|
||||
pass ``busy_timeout`` and race on WAL frames — the root cause of
|
||||
multi-writer corruption; this is defense-in-depth behind
|
||||
``_guard_supervised_gateway_conflict``. Non-blocking on purpose: the
|
||||
gateway's async watcher must never stall; the loser retries next interval.
|
||||
Without ``fcntl``/``msvcrt`` it degrades to a no-op (yields ``True``).
|
||||
"""
|
||||
lock_path = db_path.with_name(db_path.name + ".dispatch.lock")
|
||||
handle = None
|
||||
@@ -943,17 +938,11 @@ def _backfill_legacy_inflight_runs(conn: sqlite3.Connection) -> None:
|
||||
)
|
||||
|
||||
|
||||
# Legacy DBs defined these tables with a ``TEXT PRIMARY KEY`` id (or a nullable
|
||||
# ``TEXT last_event_id`` for ``kanban_notify_subs``); the current schema uses
|
||||
# ``INTEGER PRIMARY KEY AUTOINCREMENT`` / ``INTEGER NOT NULL DEFAULT 0``.
|
||||
# ``CREATE TABLE IF NOT EXISTS`` skips existing tables and
|
||||
# ``_add_column_if_missing`` only adds columns, so a drifted column type
|
||||
# requires a rebuild.
|
||||
#
|
||||
# Each entry pairs the canonical CREATE TABLE with the CREATE INDEX statements
|
||||
# DROP TABLE would take down with it (incl. ``idx_events_run`` from the additive
|
||||
# pass). ``test_rebuilt_schema_matches_fresh`` guards this list against
|
||||
# drifting from SCHEMA_SQL.
|
||||
# Legacy DBs used a ``TEXT PRIMARY KEY`` id (nullable ``TEXT last_event_id``
|
||||
# for ``kanban_notify_subs``); the additive migrations can't change a column
|
||||
# type, so drift requires a rebuild. Each entry pairs the canonical CREATE
|
||||
# TABLE with the indexes DROP TABLE takes down with it.
|
||||
# ``test_rebuilt_schema_matches_fresh`` guards this against SCHEMA_SQL drift.
|
||||
_REBUILD_SPECS = {
|
||||
"task_events": (
|
||||
"CREATE TABLE task_events ("
|
||||
@@ -1018,15 +1007,12 @@ def _table_has_drifted(conn: sqlite3.Connection, table: str) -> bool:
|
||||
def _rebuild_drifted_tables(conn: sqlite3.Connection) -> None:
|
||||
"""Rebuild any kanban table whose column types drifted from SCHEMA_SQL.
|
||||
|
||||
Old boards crash the gateway notifier (``int(None)`` on a NULL id in
|
||||
``unseen_events_for_sub``) and never match ``id > cursor``, so every
|
||||
notification is silently lost. Rebuild is CREATE new → INSERT shared
|
||||
columns → DROP old → RENAME, recreating indexes (DROP TABLE takes them
|
||||
down). Legacy TEXT ids are dropped; AUTOINCREMENT assigns fresh ones and
|
||||
``last_event_id`` cursors reset to 0, so the first post-migration tick
|
||||
replays a task's event history once — the safe failure mode for a feature
|
||||
that was already fully broken. One transaction so an interruption can't
|
||||
leave a table half-renamed, under ``connect()``'s init locks. Idempotent.
|
||||
Drifted boards crash the gateway notifier (``int(None)`` on a NULL id) and
|
||||
never match ``id > cursor``, silently losing every notification. Legacy
|
||||
TEXT ids are dropped (AUTOINCREMENT reassigns) and cursors reset to 0, so
|
||||
the first post-migration tick replays history once — safe for a feature
|
||||
that was already fully broken. One transaction under ``connect()``'s init
|
||||
locks so an interruption can't leave a table half-renamed. Idempotent.
|
||||
"""
|
||||
drifted = [t for t in _REBUILD_SPECS if _table_has_drifted(conn, t)]
|
||||
if not drifted:
|
||||
@@ -1093,13 +1079,10 @@ def _check_file_length_invariant(conn: sqlite3.Connection) -> None:
|
||||
|
||||
|
||||
# SQLite's busy_timeout backoff is near-deterministic, so stampeding writers
|
||||
# re-collide in lockstep; a jittered retry on the transaction boundary breaks
|
||||
# the convoy (mirrors state.db's _execute_write: 20-150ms band, the 20ms floor
|
||||
# prevents busy-spinning back into the collision). Only BEGIN IMMEDIATE and
|
||||
# COMMIT are retried — idempotent re-issues that touch no transaction body, so
|
||||
# a CAS inside write_txn is never replayed. Fewer retries than state.db (5 vs
|
||||
# 15): the 120s busy_timeout absorbs most waits; this is the backstop for the
|
||||
# tail where SQLite returns BUSY immediately.
|
||||
# re-collide in lockstep; a jittered 20-150ms retry on the transaction boundary
|
||||
# breaks the convoy (mirrors state.db). Only BEGIN IMMEDIATE and COMMIT are
|
||||
# retried — idempotent re-issues, so a CAS inside write_txn is never replayed.
|
||||
# 5 retries (not state.db's 15): the 120s busy_timeout absorbs most waits.
|
||||
_BUSY_MAX_RETRIES = 5
|
||||
_BUSY_RETRY_MIN_S = 0.020 # 20ms
|
||||
_BUSY_RETRY_MAX_S = 0.150 # 150ms
|
||||
@@ -1125,18 +1108,14 @@ def _execute_boundary_with_retry(conn: sqlite3.Connection, sql: str) -> None:
|
||||
|
||||
@contextlib.contextmanager
|
||||
def write_txn(conn: sqlite3.Connection, *, allow_nested: bool = False):
|
||||
"""Context manager for an IMMEDIATE write transaction; a claim CAS inside
|
||||
is atomic — at most one concurrent writer succeeds.
|
||||
"""IMMEDIATE write transaction; a claim CAS inside is atomic — at most one
|
||||
concurrent writer succeeds.
|
||||
|
||||
Nesting is an explicit opt-in (``allow_nested=True`` → savepoint instead
|
||||
of a second ``BEGIN IMMEDIATE``; otherwise a loud ``RuntimeError``). Only
|
||||
composition primitives graph builders run under one outer commit
|
||||
(``create_task``, ``add_comment``) opt in — helpers with post-commit side
|
||||
effects (``complete_task`` & co.) must never run under an open outer
|
||||
transaction, because those side effects (workspace cleanup, ready
|
||||
recomputation, failure-counter clears) would fire while the outer
|
||||
transaction can still roll back. The ROLLBACK on exception is wrapped so a
|
||||
SQLite auto-rollback doesn't shadow the original exception.
|
||||
Nesting is an explicit opt-in (``allow_nested=True`` → savepoint; otherwise
|
||||
a loud ``RuntimeError``). Only composition primitives (``create_task``,
|
||||
``add_comment``) opt in — helpers with post-commit side effects
|
||||
(``complete_task`` & co.) must never run under an open outer transaction,
|
||||
since those side effects would fire while the outer txn can still roll back.
|
||||
"""
|
||||
_kb._assert_not_delegated_child_mutation()
|
||||
if getattr(conn, "in_transaction", False):
|
||||
|
||||
@@ -154,16 +154,11 @@ def _record_worker_exit(pid: int, raw_status: int) -> None:
|
||||
|
||||
|
||||
def _classify_worker_exit(pid: int) -> "tuple[str, Optional[int]]":
|
||||
"""Return ``(kind, code)`` for a reaped worker PID.
|
||||
|
||||
``clean_exit`` (status 0 while still ``running`` = protocol violation:
|
||||
no ``kanban_complete``/``kanban_block``, so auto-block — retrying just
|
||||
loops); ``rate_limited`` (``KANBAN_RATE_LIMIT_EXIT_CODE`` — released to
|
||||
``ready`` WITHOUT counting a failure so a long quota window can't trip the
|
||||
breaker); ``nonzero_exit``; ``signaled`` (OOM/SIGKILL; ``code`` is the
|
||||
signal); ``unknown`` (pid not in the reap registry — fall back to the
|
||||
crashed-counter path; ``code`` is None).
|
||||
"""
|
||||
"""``(kind, code)`` for a reaped worker PID: ``clean_exit`` (rc 0 while
|
||||
still ``running`` = protocol violation), ``rate_limited``
|
||||
(``KANBAN_RATE_LIMIT_EXIT_CODE``, never counts as a failure),
|
||||
``nonzero_exit``, ``signaled`` (``code`` is the signal), ``unknown`` (pid
|
||||
not in the reap registry; ``code`` None)."""
|
||||
entry = _recent_worker_exits.get(int(pid))
|
||||
if entry is None:
|
||||
return ("unknown", None)
|
||||
@@ -387,18 +382,12 @@ def heartbeat_worker(
|
||||
"""
|
||||
now = int(time.time())
|
||||
with _kb.write_txn(conn):
|
||||
if expected_run_id is None:
|
||||
cur = conn.execute(
|
||||
"UPDATE tasks SET last_heartbeat_at = ? "
|
||||
"WHERE id = ? AND status = 'running'",
|
||||
(now, task_id),
|
||||
)
|
||||
else:
|
||||
cur = conn.execute(
|
||||
"UPDATE tasks SET last_heartbeat_at = ? "
|
||||
"WHERE id = ? AND status = 'running' AND current_run_id = ?",
|
||||
(now, task_id, int(expected_run_id)),
|
||||
)
|
||||
sql = "UPDATE tasks SET last_heartbeat_at = ? WHERE id = ? AND status = 'running'"
|
||||
params: tuple = (now, task_id)
|
||||
if expected_run_id is not None:
|
||||
sql += " AND current_run_id = ?"
|
||||
params += (int(expected_run_id),)
|
||||
cur = conn.execute(sql, params)
|
||||
if cur.rowcount != 1:
|
||||
return False
|
||||
run_id = (
|
||||
@@ -517,15 +506,12 @@ def detect_stale_running(
|
||||
"""Reclaim ``running`` tasks with no heartbeat progress; returns their ids.
|
||||
|
||||
Stale = running longer than ``stale_timeout_seconds`` (active run's
|
||||
``started_at``, else ``tasks.started_at``) AND ``last_heartbeat_at`` older
|
||||
than ``_STALE_HEARTBEAT_GAP_SECONDS`` or NULL. Task returns to its source
|
||||
phase, run closes ``outcome='stale'``, live host-local worker is terminated.
|
||||
``stale_timeout_seconds=0`` disables the check. ``signal_fn`` is a test hook.
|
||||
|
||||
Deliberately does NOT call ``_record_task_failure``: stale reclaim is
|
||||
dispatcher-side detection of an absent heartbeat, not a worker failure;
|
||||
counting it would let legitimately long-running tasks (>4h, no explicit
|
||||
heartbeat) trip the breaker. The ``stale`` event is the audit surface.
|
||||
``started_at``, else ``tasks.started_at``) AND ``last_heartbeat_at`` NULL or
|
||||
older than ``_STALE_HEARTBEAT_GAP_SECONDS``. Task returns to its source
|
||||
phase, run closes ``outcome='stale'``, a live host-local worker is killed.
|
||||
``0`` disables the check; ``signal_fn`` is a test hook. Deliberately NOT
|
||||
counted via ``_record_task_failure``: an absent heartbeat is not a worker
|
||||
failure, and counting it would let long-running tasks trip the breaker.
|
||||
"""
|
||||
if stale_timeout_seconds <= 0:
|
||||
return []
|
||||
@@ -938,18 +924,12 @@ def _account_crashes(conn: sqlite3.Connection, crash_details: list) -> list[str]
|
||||
def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]:
|
||||
"""Reclaim ``running`` tasks whose worker PID is no longer alive.
|
||||
|
||||
Appends ``crashed`` and restores the task's source phase; unlike
|
||||
``release_stale_claims`` it checks liveness immediately rather than waiting
|
||||
for the claim TTL. Only tasks claimed by *this host* — PIDs from other hosts
|
||||
are meaningless, and ``_default_spawn`` always runs workers locally.
|
||||
|
||||
Clean exit (rc=0) while still ``running`` is a protocol violation (no
|
||||
``kanban_complete``/``kanban_block``); it gets a bounded violation-only retry
|
||||
budget (``_protocol_violation_streak``). ``KANBAN_RATE_LIMIT_EXIT_CODE`` is a
|
||||
quota wall, NOT a failure: released WITHOUT counting a failure and stamped
|
||||
with a quota-blocker error so ``check_respawn_guard`` defers respawn. Those
|
||||
ids surface via the ``_last_rate_limited`` function attribute; the return
|
||||
stays crashed-only.
|
||||
Restores the source phase immediately (no waiting for the claim TTL), for
|
||||
tasks claimed by *this host* only — other hosts' PIDs are meaningless.
|
||||
Clean exit while ``running`` is a protocol violation with a bounded
|
||||
violation-only retry budget; ``KANBAN_RATE_LIMIT_EXIT_CODE`` is a quota
|
||||
wall, released WITHOUT counting a failure and surfaced via the
|
||||
``_last_rate_limited`` attribute (the return stays crashed-only).
|
||||
"""
|
||||
sweep = _reclaim_dead_workers(conn)
|
||||
# Outside the main txn: account each crash and maybe trip the breaker.
|
||||
@@ -986,20 +966,17 @@ def _record_task_failure(
|
||||
end_run: bool = False,
|
||||
event_payload_extra: Optional[dict] = None,
|
||||
) -> bool:
|
||||
"""Record a non-success outcome (spawn_failed / crashed / timed_out) and
|
||||
maybe trip the circuit breaker. Every non-success path funnels through
|
||||
here so ``consecutive_failures`` stays consistent. Returns True when the
|
||||
task was auto-blocked.
|
||||
"""Record a non-success outcome and maybe trip the circuit breaker; every
|
||||
non-success path funnels through here so ``consecutive_failures`` stays
|
||||
consistent. Returns True when the task was auto-blocked.
|
||||
|
||||
``release_claim=True, end_run=True`` is the spawn-failure path (task still
|
||||
running with an open run: restore source phase — or ``blocked`` on trip —
|
||||
release the claim, close the run). ``release_claim=False, end_run=False``
|
||||
is the timeout/crash path (caller ALREADY restored the phase and closed the
|
||||
run; only the counter moves, a trip re-transitions to ``blocked`` with
|
||||
``gave_up``). Threshold: per-task ``max_retries`` (nothing overrides it),
|
||||
then ``failure_limit``, then ``DEFAULT_FAILURE_LIMIT``. ``force_trip=True``
|
||||
trips unconditionally — the caller already applied its own bounded-retry
|
||||
policy; the failure is still counted.
|
||||
``release_claim=True, end_run=True``: spawn-failure path (task still
|
||||
running with an open run — restore source phase or ``blocked``, release
|
||||
claim, close run). Both False: timeout/crash path (caller already restored
|
||||
the phase and closed the run; only the counter moves, a trip flips to
|
||||
``blocked`` + ``gave_up``). Threshold: per-task ``max_retries`` >
|
||||
``failure_limit`` > ``DEFAULT_FAILURE_LIMIT``. ``force_trip`` trips
|
||||
unconditionally (caller applied its own bounded-retry policy).
|
||||
"""
|
||||
if failure_limit is None:
|
||||
failure_limit = DEFAULT_FAILURE_LIMIT
|
||||
@@ -1138,31 +1115,17 @@ def check_respawn_guard(
|
||||
) -> Optional[str]:
|
||||
"""Return a guard reason if ``task_id`` should NOT be re-spawned, else None.
|
||||
|
||||
Called per ready/review task in ``dispatch_once`` before any claim attempt;
|
||||
a reason defers the spawn this tick. The review lane skips
|
||||
``recent_success`` and ``active_pr``: a recent PR comment / completed run
|
||||
is the *precondition* of a review handoff, not a duplicate-work signal.
|
||||
|
||||
Checks in priority order:
|
||||
|
||||
``"rate_limit_cooldown"`` — latest run ended ``rate_limited`` within
|
||||
``_resolve_rate_limit_cooldown_seconds()``. Checked BEFORE
|
||||
``blocker_auth`` because the rate-limit requeue stamps a quota-flavored
|
||||
``last_failure_error`` that would otherwise match the auth regex and
|
||||
park the task forever (that path never increments
|
||||
``consecutive_failures``, so the breaker can't free it).
|
||||
``"blocker_auth"`` — ``last_failure_error`` matches a quota/auth pattern.
|
||||
``consecutive_failures`` still trips the breaker eventually, so a
|
||||
persistent auth error blocks while a transient 429 gets a few ticks.
|
||||
``"recent_success"`` — completed run within ``_RESPAWN_GUARD_SUCCESS_WINDOW``.
|
||||
Bypassed when a re-queue event (status, promote, unblock, reclaim)
|
||||
arrives AFTER that completion — a deliberate re-run.
|
||||
``"active_pr"`` — GitHub PR URL in a comment within
|
||||
``_RESPAWN_GUARD_PR_WINDOW``; re-spawning risks a duplicate PR.
|
||||
|
||||
Stale / dead claim locks are NOT a guard reason — ``release_stale_claims``
|
||||
and ``detect_crashed_workers`` reset those only after verifying the lock
|
||||
is genuinely dead.
|
||||
Called per ready/review row before any claim attempt. Priority order:
|
||||
``"rate_limit_cooldown"`` (latest run ``rate_limited`` within the cooldown;
|
||||
checked BEFORE ``blocker_auth`` because the requeue stamps a quota-flavored
|
||||
``last_failure_error`` that would otherwise park the task forever — that
|
||||
path never increments ``consecutive_failures``), ``"blocker_auth"``
|
||||
(quota/auth pattern; the breaker still trips eventually), then for the
|
||||
ready lane only ``"recent_success"`` (completed run within the window, unless
|
||||
a re-queue event arrived after it — a deliberate re-run) and ``"active_pr"``
|
||||
(PR URL in a recent comment; re-spawning risks a duplicate PR). The review
|
||||
lane skips the last two: they are the *inputs* to a review handoff. Stale /
|
||||
dead claim locks are NOT a guard reason — the reclaim passes own those.
|
||||
"""
|
||||
row = conn.execute(
|
||||
"SELECT last_failure_error FROM tasks WHERE id = ?",
|
||||
@@ -1292,16 +1255,11 @@ def review_dispatch_enabled() -> bool:
|
||||
return True
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Memory-aware dispatch guard
|
||||
#
|
||||
# With no ``kanban.max_in_progress`` on a busy board the dispatcher fanned out
|
||||
# ~30 workers on a 1 GiB VM and OOM'd the host. Two safeguards: a memory-DERIVED
|
||||
# default cap when the operator never set one (``resolve_max_in_progress``), and
|
||||
# a live memory-PRESSURE guard inside the tick (``_memory_pressure_level``)
|
||||
# because a static cap can't see other tenants. Both fail open: non-Linux or
|
||||
# any read error → empty sample → no cap / "unknown".
|
||||
# ---------------------------------------------------------------------------
|
||||
# Memory-aware dispatch guard: an uncapped board once OOM'd a 1 GiB host. Two
|
||||
# safeguards — a memory-DERIVED default cap when none is configured
|
||||
# (``resolve_max_in_progress``) and a live memory-PRESSURE guard inside the
|
||||
# tick (``_memory_pressure_level``) because a static cap can't see other
|
||||
# tenants. Both fail open: non-Linux / read error → no cap / "unknown".
|
||||
|
||||
# Assumed per-worker footprint for the derived cap; deliberately conservative
|
||||
# so the cap errs toward fewer workers on small VMs.
|
||||
@@ -1777,14 +1735,11 @@ def _dispatch_once_locked(
|
||||
max_in_progress_per_profile: Optional[int] = None,
|
||||
reconcile_orphans: bool = True,
|
||||
) -> DispatchResult:
|
||||
"""Run one dispatcher tick: reclaim stale (TTL / heartbeat) and crashed
|
||||
(dead host-local PID) running tasks, promote todo -> ready where all
|
||||
parents are done, then for each ready task with an assignee atomically
|
||||
claim and call ``spawn_fn(task, workspace_path, board) -> Optional[int]``,
|
||||
recording the PID as ``worker_pid`` so later ticks catch crashes before
|
||||
the TTL expires. After ``failure_limit`` consecutive failures the task is
|
||||
auto-blocked. Cap semantics: :func:`_tick_spawn_budget`.
|
||||
"""
|
||||
"""One dispatcher tick: reclaim stale/crashed running tasks, promote
|
||||
todo -> ready, then atomically claim each spawnable ready/review row and
|
||||
call ``spawn_fn(task, workspace_path, board) -> Optional[int]``, recording
|
||||
the PID so later ticks catch crashes before the TTL. Cap semantics:
|
||||
:func:`_tick_spawn_budget`."""
|
||||
result = DispatchResult()
|
||||
_run_reclaim_phase(
|
||||
conn, result, stale_timeout_seconds=stale_timeout_seconds,
|
||||
@@ -1985,16 +1940,12 @@ def _hermes_path_argv(path: str) -> list[str]:
|
||||
|
||||
|
||||
def _resolve_hermes_argv() -> list[str]:
|
||||
"""Resolve the ``hermes`` invocation as argv parts for ``Popen``.
|
||||
|
||||
Tries in order: ``$HERMES_BIN`` (path-like values normalized to absolute;
|
||||
bare names keep PATH semantics, never a same-directory file first);
|
||||
``shutil.which("hermes")`` normalized to absolute (on Windows ``which`` can
|
||||
return a relative ``.\\hermes.CMD`` when cwd is on PATH, and batch shims are
|
||||
unsafe with task-derived argv, so those fall back to the module form);
|
||||
``sys.executable -m hermes_cli.main`` for setups where the shim is not on
|
||||
the dispatcher's ``$PATH`` (cron, systemd ``User=`` services, launchd).
|
||||
Mirrors ``gateway.run._resolve_hermes_bin``; kept local because
|
||||
"""Resolve the ``hermes`` invocation as argv for ``Popen``: ``$HERMES_BIN``
|
||||
(path-like -> absolute; bare names keep PATH semantics, never a
|
||||
same-directory file), then ``which("hermes")`` (Windows: safe PATH search,
|
||||
batch shims fall back to the module form), then ``sys.executable -m
|
||||
hermes_cli.main`` for shim-less environments (cron, systemd ``User=``,
|
||||
launchd). Mirrors ``gateway.run._resolve_hermes_bin``; local because
|
||||
``hermes_cli`` sits below ``gateway`` in the dependency order.
|
||||
"""
|
||||
import shutil
|
||||
|
||||
@@ -79,20 +79,15 @@ def add_notify_sub(
|
||||
delivery_metadata: Optional[Mapping[str, Any]] = None,
|
||||
) -> None:
|
||||
"""Register a gateway source wanting terminal-state notifications for
|
||||
``task_id``. Idempotent on (task, platform, chat, thread).
|
||||
``task_id``; idempotent on (task, platform, chat, thread).
|
||||
|
||||
``user_id_alt`` (Signal UUID, Feishu union_id, ...) must be replayed on
|
||||
active wake: ``build_session_key`` prefers it over ``user_id``, so replaying
|
||||
only ``user_id`` would key the wake into a different session when they
|
||||
diverge. ``chat_type`` is likewise replayed so the woken turn resolves the
|
||||
operator's real channel; ``None`` keeps an existing row's value.
|
||||
|
||||
``delivery_mode`` (``_NOTIFY_DELIVERY_MODES``): ``None`` leaves an existing
|
||||
row untouched (fresh rows get ``"notify"``); an explicit value is
|
||||
last-write-wins so re-subscribing can change the mode; unknown values fall
|
||||
back to ``"notify"``. New subs start "caught up" (``last_event_id`` =
|
||||
current ``MAX(task_events.id)``, not 0) — otherwise the notifier replays
|
||||
every historical terminal event on its next tick (boot-time bursts).
|
||||
``user_id_alt`` (Signal UUID, Feishu union_id, ...) and ``chat_type`` are
|
||||
replayed on active wake: ``build_session_key`` prefers the alt id, so
|
||||
omitting it would key the wake into a different session. ``None`` keeps an
|
||||
existing row's value. ``delivery_mode``: ``None`` leaves an existing row
|
||||
untouched, an explicit valid value is last-write-wins, unknown falls back
|
||||
to ``"notify"``. New subs start caught up (``last_event_id`` =
|
||||
``MAX(task_events.id)``) so the notifier never replays history at boot.
|
||||
"""
|
||||
valid_mode = delivery_mode if delivery_mode in _NOTIFY_DELIVERY_MODES else None
|
||||
# api_server is stateless: the adapter has no send(), the wake self-post IS
|
||||
@@ -202,17 +197,13 @@ def count_notify_subs(
|
||||
chat_id: Optional[str] = None,
|
||||
thread_id: Optional[str] = None,
|
||||
) -> int:
|
||||
"""Count ``kanban_notify_subs`` rows via a read-only connection.
|
||||
|
||||
Cheap probe for the notifier's zero-subscription early exit: unlike
|
||||
:func:`connect` it never creates the DB file, runs schema init/migration,
|
||||
or opens writable (a read-only WAL open may still create ``-shm``/``-wal``
|
||||
sidecars but cannot write table content). Rows in a not-yet-checkpointed
|
||||
WAL are visible, so a fresh subscription is never missed. A missing DB or
|
||||
a legacy DB without the table counts as zero. Platform matching is
|
||||
case-insensitive (matching notifier routing); chat/thread are exact.
|
||||
Path resolution matches :func:`connect`. Raises :class:`sqlite3.Error`
|
||||
when the DB exists but cannot be read; callers choose their own fallback.
|
||||
"""Count ``kanban_notify_subs`` rows via a read-only connection — the
|
||||
notifier's cheap zero-subscription early exit. Unlike :func:`connect` it
|
||||
never creates the file, runs init/migration or opens writable; WAL rows are
|
||||
still visible so a fresh sub is never missed. Missing DB / missing table
|
||||
counts as zero; platform matches case-insensitively (as notifier routing),
|
||||
chat/thread exactly. Raises :class:`sqlite3.Error` if the DB exists but is
|
||||
unreadable — callers pick their own fallback.
|
||||
"""
|
||||
path = db_path if db_path is not None else _kb.kanban_db_path(board=board)
|
||||
if not path.exists():
|
||||
@@ -266,17 +257,14 @@ def remove_notify_sub(
|
||||
|
||||
|
||||
def purge_stale_done_notify_subs(conn: sqlite3.Connection, *, max_age_days: int = 30) -> int:
|
||||
"""Delete notify subs whose task has sat in ``done``/``blocked`` untouched
|
||||
for longer than ``max_age_days``.
|
||||
"""Delete notify subs whose task sat in ``done``/``blocked`` untouched for
|
||||
longer than ``max_age_days`` (``<= 0`` disables); returns rows deleted.
|
||||
|
||||
Subs survive ``done`` because a completed task can be reopened and must
|
||||
still notify its origin; on never-archiving boards that would accumulate
|
||||
rows forever, each scanned every notifier tick. ``blocked`` tasks are
|
||||
reaped on the same clock — they are abandoned, not merely waiting like
|
||||
``backlog``/``ready``. Age is measured from the most recent event
|
||||
(falling back to ``completed_at`` then ``created_at``), so ANY activity —
|
||||
including a reopen — resets or exempts it. ``max_age_days <= 0`` disables
|
||||
the sweep. Returns the number of rows deleted.
|
||||
Subs survive ``done`` because a reopened task must still notify its origin,
|
||||
which accumulates forever on never-archiving boards. ``blocked`` is
|
||||
abandoned (unlike ``backlog``/``ready``) so it reaps on the same clock. Age
|
||||
= latest event, else ``completed_at``, else ``created_at`` — any activity,
|
||||
including a reopen, exempts the sub.
|
||||
"""
|
||||
try:
|
||||
days = int(max_age_days)
|
||||
|
||||
@@ -95,19 +95,15 @@ def _scratch_workspace(conn: sqlite3.Connection, task_id: str) -> Optional[Path]
|
||||
|
||||
|
||||
def _is_managed_scratch_path(p: Path) -> bool:
|
||||
"""Return True iff *p* is a strict descendant of a kanban-managed scratch root.
|
||||
|
||||
A managed root is exclusively a ``workspaces/`` directory:
|
||||
``HERMES_KANBAN_WORKSPACES_ROOT`` (dispatcher-injected worker override),
|
||||
``<kanban_home>/kanban/workspaces`` (legacy default board), or
|
||||
``<kanban_home>/kanban/boards/<slug>/workspaces`` per on-disk board.
|
||||
Strict descendancy: a path EQUAL to a root is not managed (deleting it
|
||||
would wipe every task's scratch dir), and ``<kanban_home>/kanban``,
|
||||
``.../logs`` or ``.../boards/<slug>`` are rejected because they hold
|
||||
Hermes' own DB, metadata and logs. :func:`_cleanup_workspace` uses this to
|
||||
refuse ``shutil.rmtree`` outside managed storage — a board
|
||||
``default_workdir`` pointing at a real source tree can otherwise pair with
|
||||
``workspace_kind='scratch'`` and make task completion delete user data.
|
||||
"""True iff *p* is a STRICT descendant of a kanban-managed ``workspaces/``
|
||||
root (``HERMES_KANBAN_WORKSPACES_ROOT``, ``<kanban_home>/kanban/workspaces``,
|
||||
or ``<kanban_home>/kanban/boards/<slug>/workspaces``). A path equal to a
|
||||
root is not managed (deleting it would wipe every task's scratch dir);
|
||||
``<kanban_home>/kanban``, ``.../logs`` and ``.../boards/<slug>`` hold
|
||||
Hermes' own DB and metadata. :func:`_cleanup_workspace` refuses
|
||||
``rmtree`` outside managed storage — a board ``default_workdir`` on a real
|
||||
source tree paired with ``workspace_kind='scratch'`` would otherwise make
|
||||
task completion delete user data.
|
||||
"""
|
||||
return _managed_scratch_path_info(p)[0]
|
||||
|
||||
@@ -488,19 +484,13 @@ def _resolve_worktree_workspace(task: Task, *, board: Optional[str] = None) -> t
|
||||
def resolve_workspace(task: Task, *, board: Optional[str] = None) -> Path:
|
||||
"""Resolve (and create if needed) the workspace for a task.
|
||||
|
||||
- ``scratch``: fresh dir under ``<board-root>/workspaces/<id>/`` — the
|
||||
same path for the dispatcher and every profile worker, so handoff is
|
||||
path-stable.
|
||||
- ``dir:<path>``: ``workspace_path``, created if missing. MUST be absolute
|
||||
— relative paths are rejected to prevent confused-deputy traversal where
|
||||
``../../../tmp/attacker`` resolves against the dispatcher's CWD.
|
||||
- ``worktree``: a real linked git worktree. A ``workspace_path`` naming a
|
||||
repo root is an anchor (materializes ``<repo>/.worktrees/<task-id>``);
|
||||
a concrete target path is created/reused; with none, the board's
|
||||
``default_workdir`` anchors it and raises if unset rather than guessing
|
||||
from the dispatcher's CWD. Empty ``branch_name`` -> ``wt/<task-id>``.
|
||||
|
||||
Persist the resolved path via ``set_workspace_path`` so later runs reuse it.
|
||||
``scratch``: ``<board-root>/workspaces/<id>/`` — path-stable across the
|
||||
dispatcher and every profile worker. ``dir``: ``workspace_path``, created
|
||||
if missing; MUST be absolute (relative paths would resolve against the
|
||||
dispatcher's CWD — confused-deputy traversal). ``worktree``: a linked git
|
||||
worktree; a repo-root ``workspace_path`` anchors ``<repo>/.worktrees/<id>``,
|
||||
a concrete path is created/reused, none -> the board's ``default_workdir``
|
||||
(raises if unset rather than guessing). Persist via ``set_workspace_path``.
|
||||
"""
|
||||
kind = task.workspace_kind or "scratch"
|
||||
if kind == "worktree":
|
||||
|
||||
Reference in New Issue
Block a user