diff --git a/README.md b/README.md
index 568dd0e..effb227 100644
--- a/README.md
+++ b/README.md
@@ -876,6 +876,73 @@ Worker публикует heartbeat в общей SQLite БД. Активный
}
```
+#### Динамические полосы и адаптивный интервал
+
+При `PP_CONCURRENCY=3` профиль может закрепить три логических слота за разными
+частями бесконечного конвейера. Это универсальная конфигурация: PromptPilot не
+знает названий этапов заранее и связывает их с очередями того же профиля:
+
+```json
+"scheduler": {
+ "lanes": [
+ {"id": "integration", "queues": ["merge", "review"], "borrow": true},
+ {"id": "production", "queues": ["review", "fix"], "borrow": true},
+ {"id": "intake", "queues": ["triage", "plan"], "borrow": true}
+ ]
+}
+```
+
+Каждая запись — один слот, `queues` задаёт его локальный порядок. Поэтому
+интеграционный слот сначала берёт MERGE, а при его отсутствии REVIEW;
+производственный — REVIEW, затем FIX; входной — TRIAGE, затем PLAN. Если своей
+работы нет, `borrow: true` разрешает слоту взять runnable-задачу другой полосы.
+Занятый слот не считается свободным до завершения точного запуска. Выбор и
+claim одной задачи происходят в одной SQLite-транзакции, поэтому два слота не
+получат одно вхождение. Без `scheduler` остаётся прежний порядок
+`priority → created_at → id`.
+
+Полоса выбирает только серию TRIAGE. Порядок целей *внутри* неё — recovery,
+ручной P0, затем обычный backlog — остаётся контрактом репозиторного `next` и
+его свежих gate-проверок. Так машинное расписание не может обойти recovery или
+самостоятельно назначить GitHub-цель.
+
+Очередь может также владеть своей частотой:
+
+```json
+{
+ "id": "fix",
+ "adaptive_cadence": {
+ "idle_recurrence": "30m",
+ "busy_recurrence": "15m",
+ "backlog_above": 3
+ }
+},
+{
+ "id": "triage",
+ "adaptive_cadence": {
+ "idle_recurrence": "30m",
+ "busy_recurrence": "15m",
+ "backlog_above": 0,
+ "empty_runs_before_idle": 2
+ }
+},
+{
+ "id": "plan",
+ "wake_when": {"field": "plan_candidates"},
+ "adaptive_cadence": {
+ "idle_recurrence": "4h",
+ "event_wake": true
+ }
+}
+```
+
+Полный свежий снимок включает быстрый интервал только выше порога. Для TRIAGE
+он снимается существующим счётчиком серии после двух последовательных `ПУСТО`;
+FIX возвращается к 30 минутам сразу при backlog ≤ 3. PLAN просыпается через
+`wake_when`/`wake_after_success`, а 4 часа остаются страховкой от потерянного
+события. Обычное чтение кэша расписание не меняет. Более длинный интервал не
+откладывает уже назначенный ранний запуск.
+
`priority_control` необязателен. Он использует GitHub-метки `queue:p0`…
`queue:p3` для ручного решения и `queue:auto:p0`…`queue:auto:p3` для оценки,
перенесённой с issue на PR. Ручная метка старше автоматической; без обеих
@@ -1362,6 +1429,16 @@ pp note 42 --clear # убрать
просыпающийся каждые два часа, превращает бота в будильник. Ошибки (`failed`)
шлются всегда, независимо от итога.
+Только targeted fallback после выполненного `next` дополнительно разрешает
+точную форму `ИТОГ: УСТАРЕЛО (gate-fallback: ...)`: exact target, HEAD или её
+executable-позиция изменились **до первой мутации**. Generic skill/workflow и
+иная форма `УСТАРЕЛО` классифицируются как настоящий `НЕ СМОГ`. Первый такой
+исход создаёт следующее вхождение сразу для свежего election; второй подряд
+идёт по обычному интервалу. Durable-ограничитель сбрасывается любым non-stale
+исходом. В дашборде валидные перевыборы считаются отдельно. Неоднозначный отказ,
+проблема доступа и ошибка после начатой мутации по-прежнему должны заканчиваться
+`НЕ СМОГ` или `НУЖЕН ЧЕЛОВЕК`, а не быстрым циклом перевыбора.
+
Парсится он **всегда**, даже когда мы не просили — если агент сам закончил такой
строкой, итог подхватится. Побеждает последнее совпадение: формат могли
процитировать по дороге, а вердикт — это закрывающая строка.
diff --git a/promptpilot/bot.py b/promptpilot/bot.py
index b249d89..96e8bd7 100644
--- a/promptpilot/bot.py
+++ b/promptpilot/bot.py
@@ -889,7 +889,6 @@ def _pipeline_text(data: dict) -> str:
health = data.get("health", {})
coverage = ("5 ч" if recent.get("complete")
else f"{recent.get('coverage_hours', 0):g} из 5 ч")
- delta = recent.get("backlog_delta") if recent.get("complete") else None
week_trend = (f"{week.get('backlog_delta'):+d}" if week.get("complete")
else f"— (покрытие {week.get('coverage_hours', 0):g} ч)")
month_trend = (f"{month.get('backlog_delta'):+d}" if month.get("complete")
@@ -903,14 +902,29 @@ def _pipeline_text(data: dict) -> str:
f"{cache_age} сек назад" if cache_age < 60 else
f"{round(cache_age / 60)} мин назад")
backlog_total = data.get("backlog_total")
+ delta_rate = recent.get("backlog_delta_per_hour")
+ entered_rate = recent.get("entered_per_hour")
+ exited_rate = recent.get("exited_per_hour")
+ rate = lambda value: "—" if value is None else f"{value:+.2f}/ч"
+ bottleneck = next(
+ (queue for queue in data.get("queues", [])
+ if queue.get("id") == data.get("bottleneck")), None)
+ outcomes = data.get("outcomes") or {}
lines = [f"📈 {data['title']}",
f"Состояние: {health.get('label', '—')} — {health.get('reason', '—')}",
f"Снимок GitHub: {cache_state}, {cache_age_text}; чтение без GitHub API",
f"Backlog: {backlog_total if backlog_total is not None else '—'}"
- + (f" (Δ 5 ч: {delta:+d})" if delta is not None else " (история копится)"),
- f"Вход / выход / переходы: {recent.get('entered', 0)} / "
- f"{recent.get('exited', 0)} / {recent.get('transitions', 0)}; покрытие {coverage}",
+ + (f" (Δ {rate(delta_rate)})" if delta_rate is not None else " (история копится)"),
+ f"Вход / выход: {rate(entered_rate)} / {rate(exited_rate)}; "
+ f"переходов {recent.get('transitions', 0)}; покрытие {coverage}",
f"Тренд backlog: 7 дней {week_trend}; 30 дней {month_trend}",
+ f"Узкое место: " + (
+ f"{bottleneck['title']} (ETA {bottleneck.get('eta_hours', '—')} ч)"
+ if bottleneck else "нет"),
+ "Безопасные ожидания / ошибки: "
+ f"{outcomes.get('safe_deferrals_now', 0) + outcomes.get('stale_reselections_5h', 0)} / "
+ f"{outcomes.get('real_errors_5h', 0)} "
+ f"(stale-перевыборов {outcomes.get('stale_reselections_5h', 0)}).",
f"Цель оценки: текущая очередь примерно за {data['target_clear_hours']:g} ч.", ""]
runtime = data.get("runtime", {})
if runtime:
@@ -969,16 +983,19 @@ def _pipeline_text(data: dict) -> str:
route_text += f" → {route.get('effective')}"
backlog = queue.get("backlog")
runs_needed = queue.get("runs_needed")
+ cadence = queue.get("adaptive_cadence") or {}
+ cadence_text = (f"; adaptive {cadence.get('mode')}"
+ if cadence else "")
lines.append(f"{marker} {queue['title']}: {backlog if backlog is not None else '—'} / "
f"{queue['capacity']} за прогон = {runs_needed if runs_needed is not None else '—'} прогонов; "
- f"сейчас {queue['interval'] or 'не настроено'}; "
+ f"сейчас {queue['interval'] or 'не настроено'}{cadence_text}; "
f"средний запуск {duration_text}; ETA {eta if eta is not None else '—'} ч\n"
f" Маршрут: {route_text}\n"
f" Рекомендация: {queue['recommendation']}")
runs = recent.get("runs", {})
lines.extend(["", f"Прогоны за окно: {runs.get('runs', 0)}; готово {runs.get('ready', 0)}, "
f"нужен человек {runs.get('human', 0)}, не смог {runs.get('unable', 0)}, "
- f"упало {runs.get('failed', 0)}.",
+ f"устарело {runs.get('stale', 0)}, упало {runs.get('failed', 0)}.",
"⚠ — текущее узкое место по ETA с учётом интервала и средней длительности."])
return "\n".join(lines)
@@ -3298,6 +3315,14 @@ async def cb_windows_list(update: Update, context: ContextTypes.DEFAULT_TYPE):
# Entry point
# ---------------------------------------------------------------------------
+def _silent_completion(task) -> bool:
+ """Whether a completed routine outcome should stay visible but not ping."""
+ return (
+ task.status.value == "completed"
+ and (task.verdict or "").upper() in {"ПУСТО", "УСТАРЕЛО"}
+ )
+
+
async def _notify_loop(bot):
"""Background loop: send notifications for completed/failed tasks every 10s.
@@ -3343,10 +3368,11 @@ async def _notify_loop(bot):
continue
for task in pending:
try:
- # ПУСТО — «проснулся по расписанию, делать нечего»: рутина
- # повторяющихся задач, ради которой будить человека не за чем.
- # Итог остаётся в базе и виден в списке задач.
- if task.status.value == "completed" and (task.verdict or "").upper() == "ПУСТО":
+ # ПУСТО and a validated targeted УСТАРЕЛО are routine queue
+ # outcomes. Both remain visible in history/metrics; only the
+ # Telegram success ping is suppressed. Invalid/generic stale
+ # text is normalized to НЕ СМОГ before it reaches this point.
+ if _silent_completion(task):
db.mark_notified(task.id)
continue
if task.status.value == "completed":
diff --git a/promptpilot/db.py b/promptpilot/db.py
index ca9e622..276c2c6 100644
--- a/promptpilot/db.py
+++ b/promptpilot/db.py
@@ -632,7 +632,8 @@ def recent_working_dirs(limit: int = 8, machine: Optional[str] = None) -> list:
return [r["working_dir"] for r in rows]
-def get_next_runnable(busy_keys=(), key_fn=None) -> Optional[TaskInDB]:
+def get_next_runnable(busy_keys=(), key_fn=None,
+ order_key_fn=None) -> Optional[TaskInDB]:
"""Claim the highest-priority runnable task and mark it running.
The whole select-then-claim runs under the write lock, so several workers
@@ -641,6 +642,10 @@ def get_next_runnable(busy_keys=(), key_fn=None) -> Optional[TaskInDB]:
busy_keys/key_fn — when the caller already runs tasks, candidates whose
key_fn(task) is in busy_keys are passed over: that is how two agents are
kept out of one work tree while the queue keeps moving.
+
+ order_key_fn — optional policy rank evaluated while the same write
+ transaction owns the runnable snapshot. Returning ``None`` excludes a
+ candidate; otherwise the smallest key wins before the stable DB order.
"""
now = _now()
busy = set(busy_keys or ())
@@ -650,7 +655,7 @@ def get_next_runnable(busy_keys=(), key_fn=None) -> Optional[TaskInDB]:
# task is found — a hard LIMIT could hide a free task behind a wall of
# conflicting ones. Rows are materialised before any UPDATE so claiming
# one doesn't disturb the iteration.
- limit_clause = "" if busy else " LIMIT 1"
+ limit_clause = "" if busy or order_key_fn is not None else " LIMIT 1"
rows = conn.execute(
f"""SELECT * FROM tasks
WHERE status IN ('pending', 'rate_limited')
@@ -662,10 +667,18 @@ def get_next_runnable(busy_keys=(), key_fn=None) -> Optional[TaskInDB]:
ORDER BY priority ASC, created_at ASC, id ASC{limit_clause}""",
(now, now),
).fetchall()
- for row in rows:
+ candidates = []
+ for position, row in enumerate(rows):
task = _row_to_task(row)
if busy and key_fn and key_fn(task) in busy:
continue
+ rank = order_key_fn(task) if order_key_fn is not None else ()
+ if order_key_fn is not None and rank is None:
+ continue
+ candidates.append((rank, position, task))
+ if order_key_fn is not None:
+ candidates.sort(key=lambda item: (item[0], item[1]))
+ for _rank, _position, task in candidates:
started_at = _now()
cur = conn.execute(
"""UPDATE tasks SET status = 'running', started_at = ?,
@@ -1076,6 +1089,136 @@ def update_series(series_id: int, fields: dict) -> bool:
return True
+def apply_pipeline_series_cadence(
+ series_id: int, *, idle_recurrence: str,
+ busy_recurrence: Optional[str] = None, boost: bool = False,
+ empty_runs_before_idle: int = 0,
+ publication_guard: Optional[dict] = None) -> Optional[dict]:
+ """Idempotently reconcile one profile-owned adaptive cadence.
+
+ ``idle_recurrence`` is the durable safety interval. ``busy_recurrence``
+ uses the existing temporary-boost fields, so completion and crash recovery
+ keep one source of truth. When ``empty_runs_before_idle`` is non-zero, a
+ quiet snapshot does not switch cadence by itself: consecutive ``ПУСТО``
+ verdicts retire the boost in :func:`prepare_series_recurrence`.
+
+ A profile that opts into this API owns these recurrence fields. The write
+ is one transaction and never postpones an occurrence already scheduled
+ earlier than the new effective cadence.
+ """
+ if (not idle_recurrence or parse_recurrence(idle_recurrence) is None
+ or (busy_recurrence is not None
+ and parse_recurrence(busy_recurrence) is None)):
+ raise ValueError("invalid adaptive pipeline recurrence")
+ if (isinstance(empty_runs_before_idle, bool)
+ or not isinstance(empty_runs_before_idle, int)
+ or not 0 <= empty_runs_before_idle <= 20):
+ raise ValueError("empty_runs_before_idle must be from 0 to 20")
+
+ if publication_guard is not None:
+ required = {
+ "epoch_key", "epoch", "epoch_default", "profile_key",
+ "profile_hash", "revision_key", "revision",
+ }
+ if (not isinstance(publication_guard, dict)
+ or set(publication_guard) != required
+ or any(not isinstance(publication_guard[name], str)
+ for name in required)):
+ raise ValueError("invalid adaptive cadence publication guard")
+
+ with _connect(immediate=True) as conn:
+ if publication_guard is not None:
+ def setting(name: str, default: Optional[str] = None):
+ row = conn.execute(
+ "SELECT value FROM settings WHERE key = ?",
+ (publication_guard[name],),
+ ).fetchone()
+ return row["value"] if row else default
+
+ if (setting("epoch_key", publication_guard["epoch_default"])
+ != publication_guard["epoch"]
+ or setting("profile_key")
+ != publication_guard["profile_hash"]
+ or setting("revision_key")
+ != publication_guard["revision"]):
+ return None
+ row = conn.execute(
+ "SELECT * FROM task_series WHERE id = ?", (series_id,)
+ ).fetchone()
+ if not row or row["ended_at"]:
+ return None
+ before = dict(row)
+ desired = dict(before)
+ desired["base_recurrence"] = idle_recurrence
+
+ if busy_recurrence is not None:
+ keep_until_empty = (
+ not boost and empty_runs_before_idle > 0
+ and before.get("temporary_recurrence") == busy_recurrence
+ )
+ if boost:
+ desired["temporary_recurrence"] = busy_recurrence
+ desired["temporary_until"] = None
+ desired["temporary_empty_limit"] = (
+ empty_runs_before_idle or None)
+ # A live busy observation breaks an empty-run streak even if
+ # the effective interval itself was already boosted.
+ desired["temporary_empty_count"] = 0
+ elif not keep_until_empty:
+ desired["temporary_recurrence"] = None
+ desired["temporary_until"] = None
+ desired["temporary_empty_limit"] = None
+ desired["temporary_empty_count"] = 0
+ else:
+ # An adaptive policy owns the whole cadence. Event-only queues do
+ # not inherit an unrelated/manual temporary boost forever.
+ desired["temporary_recurrence"] = None
+ desired["temporary_until"] = None
+ desired["temporary_empty_limit"] = None
+ desired["temporary_empty_count"] = 0
+
+ changed_fields = [
+ name for name in (
+ "base_recurrence", "temporary_recurrence", "temporary_until",
+ "temporary_empty_limit", "temporary_empty_count",
+ )
+ if desired.get(name) != before.get(name)
+ ]
+ if changed_fields:
+ now = _now()
+ assignments = ", ".join(f"{name} = ?" for name in changed_fields)
+ conn.execute(
+ f"UPDATE task_series SET {assignments}, updated_at = ? WHERE id = ?",
+ (*[desired.get(name) for name in changed_fields], now, series_id),
+ )
+ if "base_recurrence" in changed_fields:
+ conn.execute(
+ """UPDATE tasks SET recurrence = ? WHERE series_id = ?
+ AND status IN ('pending', 'rate_limited')""",
+ (idle_recurrence, series_id),
+ )
+
+ effective = _effective_series_recurrence(desired)
+ candidate = parse_recurrence(effective)
+ if candidate:
+ candidate_iso = _to_utc_iso(candidate)
+ conn.execute(
+ """UPDATE tasks SET scheduled_at = ?
+ WHERE series_id = ? AND status = 'pending'
+ AND scheduled_at > ?""",
+ (candidate_iso, series_id, candidate_iso),
+ )
+
+ return {
+ "changed": bool(changed_fields),
+ "effective_recurrence": _effective_series_recurrence(desired),
+ "base_recurrence": desired["base_recurrence"],
+ "temporary_recurrence": desired.get("temporary_recurrence"),
+ "temporary_empty_limit": desired.get("temporary_empty_limit"),
+ "temporary_empty_count": desired.get("temporary_empty_count") or 0,
+ }
+
+
def _pipeline_series_wake_intent_key(series_id: int) -> str:
if isinstance(series_id, bool) or not isinstance(series_id, int) or series_id <= 0:
raise ValueError("pipeline series id must be a positive integer")
@@ -1441,13 +1584,35 @@ def wake_series_once(series_id: Optional[int], latch_key: str,
return True
+def _pipeline_stale_reselect_guard_key(series_id: int) -> str:
+ return f"pipeline_stale_reselect_guard:v1:{int(series_id)}"
+
+
def prepare_series_recurrence(series_id: int, verdict: Optional[str]) -> Optional[dict]:
"""Update temporary-boost counters and return settings for the next run."""
- with _connect() as conn:
+ # BEGIN IMMEDIATE serializes the read/modify/write with a live cadence
+ # reconciliation. A deferred transaction could read the old empty counter
+ # and then overwrite a concurrent busy-observation reset.
+ with _connect(immediate=True) as conn:
row = conn.execute("SELECT * FROM task_series WHERE id = ?", (series_id,)).fetchone()
if not row or row["ended_at"]:
return None
s = dict(row)
+ stale_key = _pipeline_stale_reselect_guard_key(series_id)
+ stale = (verdict or "").upper() == "УСТАРЕЛО"
+ stale_reselect_immediate = False
+ if stale:
+ prior = conn.execute(
+ "SELECT value FROM settings WHERE key = ?", (stale_key,)
+ ).fetchone()
+ if prior is None:
+ conn.execute(
+ "INSERT INTO settings (key, value) VALUES (?, '1')",
+ (stale_key,),
+ )
+ stale_reselect_immediate = True
+ else:
+ conn.execute("DELETE FROM settings WHERE key = ?", (stale_key,))
now = datetime.now(timezone.utc)
until = _parse_dt(s["temporary_until"])
expired = bool(until and until <= now)
@@ -1470,6 +1635,7 @@ def prepare_series_recurrence(series_id: int, verdict: Optional[str]) -> Optiona
(empty_count, _now(), series_id))
s["temporary_empty_count"] = empty_count
s["effective_recurrence"] = _effective_series_recurrence(s, now)
+ s["stale_reselect_immediate"] = stale_reselect_immediate
return s
@@ -1646,7 +1812,8 @@ def pipeline_run_metrics(series_ids: list[int], since: datetime) -> dict:
"""Semantic outcomes of scheduled runs belonging to one pipeline profile."""
ids = sorted({int(value) for value in series_ids if value is not None})
empty = {"runs": 0, "ready": 0, "empty": 0, "human": 0,
- "no_change": 0, "unable": 0, "failed": 0, "other": 0,
+ "no_change": 0, "stale": 0, "unable": 0, "failed": 0,
+ "other": 0,
"unresolved_unable": 0, "unresolved_failed": 0,
"recovered_unable": 0, "recovered_failed": 0,
"tokens_known_runs": 0, "input_tokens": 0,
@@ -1696,6 +1863,11 @@ def pipeline_run_metrics(series_ids: list[int], since: datetime) -> dict:
# The stage itself completed correctly even if the selected item
# now waits for a person, so an older execution incident is over.
unresolved[int(row["series_id"])] = {"unable": 0, "failed": 0}
+ elif verdict == "УСТАРЕЛО":
+ # An exact target changed before the first mutation. This is a
+ # successful safety fence followed by re-election, not an error.
+ result["stale"] += 1
+ unresolved[int(row["series_id"])] = {"unable": 0, "failed": 0}
elif verdict == "НЕ СМОГ":
result["unable"] += 1
unresolved[int(row["series_id"])]["unable"] += 1
@@ -1805,8 +1977,13 @@ def set_setting_if_newer_revision(
key: str, value: str, *, revision_key: str, revision: int,
guard_key: str, expected_guard: str,
guard_default: Optional[str] = None,
- lease_guard: Optional[dict] = None) -> bool:
+ lease_guard: Optional[dict] = None,
+ publication_revision_key: Optional[str] = None,
+ companion_key: Optional[str] = None,
+ companion_value: Optional[str] = None) -> bool:
"""Atomically publish a newer revision while a durable guard matches."""
+ if (companion_key is None) != (companion_value is None):
+ raise ValueError("companion_key and companion_value must be set together")
with _connect(immediate=True) as conn:
if (lease_guard is not None
and not _pipeline_scan_lease_guard_matches(conn, lease_guard)):
@@ -1826,6 +2003,17 @@ def set_setting_if_newer_revision(
published_revision = 0
if published_revision >= int(revision):
return False
+ if publication_revision_key is not None:
+ row = conn.execute(
+ "SELECT value FROM settings WHERE key = ?",
+ (publication_revision_key,),
+ ).fetchone()
+ try:
+ active_revision = int(row["value"]) if row else 0
+ except (TypeError, ValueError):
+ active_revision = 0
+ if active_revision >= int(revision):
+ return False
conn.execute(
"INSERT OR REPLACE INTO settings (key, value) VALUES (?, ?)",
(key, value),
@@ -1834,6 +2022,16 @@ def set_setting_if_newer_revision(
"INSERT OR REPLACE INTO settings (key, value) VALUES (?, ?)",
(revision_key, str(int(revision))),
)
+ if companion_key is not None:
+ conn.execute(
+ "INSERT OR REPLACE INTO settings (key, value) VALUES (?, ?)",
+ (companion_key, companion_value),
+ )
+ if publication_revision_key is not None:
+ conn.execute(
+ "INSERT OR REPLACE INTO settings (key, value) VALUES (?, ?)",
+ (publication_revision_key, str(int(revision))),
+ )
return True
diff --git a/promptpilot/herdr_exec.py b/promptpilot/herdr_exec.py
index b69780f..9734331 100644
--- a/promptpilot/herdr_exec.py
+++ b/promptpilot/herdr_exec.py
@@ -77,6 +77,10 @@
r"(?:\s*(?:[—-]\s*.*|\([^\r\n)]*\)))?$",
re.IGNORECASE,
)
+TARGETED_STALE_CLOSING_VERDICT_RE = re.compile(
+ r"^ИТОГ:\s*УСТАРЕЛО\s+\(gate-fallback:\s*(?=\S).+\)$",
+ re.IGNORECASE,
+)
AGY_BACKGROUND_RUNNING_RE = re.compile(
r"(?mi)^(?:\s*[●•]\s*\[[^\]]+\].*\brunning\s*"
r"|\s*[\u2800-\u28ff]\s+Running command(?:\.\.\.)?\s*)$"
@@ -87,12 +91,20 @@ class HerdrError(Exception):
pass
-def ensure_closing_verdict_contract(prompt: str) -> str:
+def ensure_closing_verdict_contract(
+ prompt: str, *, allow_targeted_stale: bool = False) -> str:
"""Add a trusted response boundary and closing-verdict contract once."""
if (WORKFLOW_CONTRACT_MARKER in prompt
and prompt.rstrip().endswith(WORKFLOW_CONTRACT_END)):
return prompt
- return f"{prompt.rstrip()}\n\n{WORKFLOW_CONTRACT_SUFFIX}"
+ suffix = WORKFLOW_CONTRACT_SUFFIX
+ if allow_targeted_stale:
+ suffix = suffix.replace(
+ "ИТОГ: ПУСТО (краткая причина)",
+ "ИТОГ: УСТАРЕЛО (gate-fallback: точная причина)\n"
+ "ИТОГ: ПУСТО (краткая причина)",
+ )
+ return f"{prompt.rstrip()}\n\n{suffix}"
def herdr_argv(args, host=None) -> list:
@@ -411,7 +423,8 @@ def _looks_env_failure(cleaned: str) -> str:
return env_failure(response_only)
-def _closing_workflow_verdict(cleaned: str) -> str:
+def _closing_workflow_verdict(
+ cleaned: str, *, allow_targeted_stale: bool = False) -> str:
"""Return a workflow verdict only when it is the final response line.
The prompt itself can contain all allowed verdict examples. Searching the
@@ -435,9 +448,17 @@ def _closing_workflow_verdict(cleaned: str) -> str:
))
break
for candidate in candidates:
+ if (allow_targeted_stale
+ and TARGETED_STALE_CLOSING_VERDICT_RE.fullmatch(candidate)):
+ return "УСТАРЕЛО"
match = WORKFLOW_CLOSING_VERDICT_RE.fullmatch(candidate)
if match:
return match.group(1).upper()
+ if re.match(r"^ИТОГ:\s*УСТАРЕЛО\b", candidate, re.IGNORECASE):
+ # Only a targeted fallback route may safely ask for re-election,
+ # and only with its exact gate-fallback evidence wrapper. Treat an
+ # invented or malformed stale result as a real execution failure.
+ return "НЕ СМОГ"
return ""
@@ -448,7 +469,8 @@ def _has_running_background_task(text: str) -> bool:
def _stabilize_workflow_completion(name, prompt, state, raw, deadline,
cancel_check, on_blocked=None, host=None,
- require_closing_verdict=False):
+ require_closing_verdict=False,
+ allow_targeted_stale=False):
"""Do not treat a transient idle between agy background tasks as done.
Antigravity can briefly return to an idle prompt while a managed background
@@ -481,7 +503,8 @@ def _stabilize_workflow_completion(name, prompt, state, raw, deadline,
background_running = _has_running_background_task(recent)
if (state in {"idle", "done"} and not background_running
and _closing_workflow_verdict(
- _trim_transcript(recent, prompt))):
+ _trim_transcript(recent, prompt),
+ allow_targeted_stale=allow_targeted_stale)):
return state, raw
rc, data, status_raw = _run(["agent", "get", name], host=host)
@@ -542,7 +565,8 @@ def run_in_herdr(task, provider_cfg: dict, on_blocked=None, timeout: int = None,
cancel_check=None, keep_pane: bool = None, host: str = None,
on_worktree=None, on_pane=None, on_started=None,
prompt_override: str = None,
- require_closing_verdict: bool = False) -> dict:
+ require_closing_verdict: bool = False,
+ allow_targeted_stale: bool = False) -> dict:
"""Run a task in a herdr-managed agent session.
on_blocked(pane_id) is called once when the agent first enters ``blocked``.
@@ -576,7 +600,8 @@ def run_in_herdr(task, provider_cfg: dict, on_blocked=None, timeout: int = None,
require_closing_verdict or WORKFLOW_CONTRACT_MARKER in prompt
)
if require_closing_verdict:
- prompt = ensure_closing_verdict_contract(prompt)
+ prompt = ensure_closing_verdict_contract(
+ prompt, allow_targeted_stale=allow_targeted_stale)
if requires_verdict and getattr(task, "detached", False):
outcome["error"] = (
"herdr detached mode is incompatible with a required closing "
@@ -825,6 +850,7 @@ def fail(msg, keep_pane=True):
name, prompt, state, raw, deadline, cancel_check,
on_blocked=on_blocked, host=host,
require_closing_verdict=requires_verdict,
+ allow_targeted_stale=allow_targeted_stale,
)
if state == "__cancel__":
@@ -881,7 +907,8 @@ def fail(msg, keep_pane=True):
outcome["error"] = cleaned
return outcome
- closing_verdict = _closing_workflow_verdict(cleaned)
+ closing_verdict = _closing_workflow_verdict(
+ cleaned, allow_targeted_stale=allow_targeted_stale)
if requires_verdict and not closing_verdict:
return fail(
"herdr agent finished without the required closing ИТОГ verdict",
diff --git a/promptpilot/models.py b/promptpilot/models.py
index 17c5604..7174fb4 100644
--- a/promptpilot/models.py
+++ b/promptpilot/models.py
@@ -88,7 +88,7 @@ class TaskInDB(BaseModel):
worktree_branch: Optional[str] = None
herdr_pane: Optional[str] = None # pane of a herdr-executor run (📺 in the bot)
note: Optional[str] = None # the human's late word, injected into the next attempt
- verdict: Optional[str] = None # ГОТОВО | УЖЕ СДЕЛАНО | НУЖЕН ЧЕЛОВЕК | НЕ СМОГ | ПУСТО (тихий: без TG-уведомления)
+ verdict: Optional[str] = None # ГОТОВО | УЖЕ СДЕЛАНО | НУЖЕН ЧЕЛОВЕК | НЕ СМОГ | УСТАРЕЛО | ПУСТО
series_id: Optional[int] = None
series_title: Optional[str] = None
series_paused: bool = False
diff --git a/promptpilot/pipeline_insights.py b/promptpilot/pipeline_insights.py
index 18990c8..fee1d7d 100644
--- a/promptpilot/pipeline_insights.py
+++ b/promptpilot/pipeline_insights.py
@@ -34,6 +34,9 @@
_CACHE_KEY_PREFIX = "pipeline_insights_cache:v1:"
_CACHE_REFRESH_REVISION_PREFIX = "pipeline_insights_refresh_revision:v1:"
_CACHE_PUBLISHED_REVISION_PREFIX = "pipeline_insights_published_revision:v1:"
+_CACHE_PUBLISHED_PROFILE_PREFIX = "pipeline_insights_published_profile:v1:"
+_CACHE_PUBLISHED_PROFILE_REVISION_PREFIX = \
+ "pipeline_insights_published_profile_revision:v1:"
_INTERVAL_PRESETS = ((0.25, "15m"), (0.5, "30m"), (1, "1h"), (2, "2h"),
(4, "4h"), (8, "8h"), (12, "12h"), (24, "24h"))
_HISTORY_WINDOWS = (5, 24, 24 * 7, 24 * 30)
@@ -955,6 +958,15 @@ def _published_revision_key(profile_id: str, profile_hash: str) -> str:
f"{quote(profile_id, safe='')}:{profile_hash}")
+def _published_profile_key(profile_id: str) -> str:
+ return f"{_CACHE_PUBLISHED_PROFILE_PREFIX}{quote(profile_id, safe='')}"
+
+
+def _published_profile_revision_key(profile_id: str) -> str:
+ return (f"{_CACHE_PUBLISHED_PROFILE_REVISION_PREFIX}"
+ f"{quote(profile_id, safe='')}")
+
+
def _profile_fingerprint(profile: dict) -> str:
# Admission tuning changes when a scan may run, not what the observation
# means. Keep last-good data readable when an operator enables/tunes the
@@ -1075,7 +1087,11 @@ def _publish_cache(profile_id: str, profile: dict, generation: int, epoch: int,
revision=revision,
guard_key=_CACHE_EPOCH_KEY, expected_guard=str(epoch),
guard_default="0",
- lease_guard=lease.guard if lease is not None else None):
+ lease_guard=lease.guard if lease is not None else None,
+ publication_revision_key=(
+ _published_profile_revision_key(profile_id)),
+ companion_key=_published_profile_key(profile_id),
+ companion_value=profile_hash):
if lease is not None:
# The CAS can also lose to an ordinary cache invalidation or a
# newer refresh. Distinguish that benign race from a rejected
@@ -1355,6 +1371,109 @@ def _interval_hours(value: str | None) -> float | None:
return None
+def _adaptive_cadence_policy(queue: dict) -> dict | None:
+ """Validate the generic queue-owned cadence policy, if configured."""
+ raw = queue.get("adaptive_cadence")
+ if raw is None:
+ return None
+ if not isinstance(raw, dict):
+ raise ValueError("adaptive_cadence должен быть JSON-объектом")
+ allowed = {
+ "idle_recurrence", "busy_recurrence", "backlog_above",
+ "empty_runs_before_idle", "event_wake",
+ }
+ unknown = sorted(set(raw) - allowed)
+ if unknown:
+ raise ValueError(
+ "adaptive_cadence содержит неизвестные поля: " + ", ".join(unknown))
+ idle = raw.get("idle_recurrence")
+ busy = raw.get("busy_recurrence")
+ if not isinstance(idle, str) or db.parse_recurrence(idle) is None:
+ raise ValueError("adaptive_cadence.idle_recurrence не разобран")
+ if busy is not None and (
+ not isinstance(busy, str) or db.parse_recurrence(busy) is None):
+ raise ValueError("adaptive_cadence.busy_recurrence не разобран")
+ threshold = raw.get("backlog_above", 0)
+ empty_runs = raw.get("empty_runs_before_idle", 0)
+ event_wake = raw.get("event_wake", False)
+ if isinstance(threshold, bool) or not isinstance(threshold, int) or threshold < 0:
+ raise ValueError("adaptive_cadence.backlog_above должен быть >= 0")
+ if (isinstance(empty_runs, bool) or not isinstance(empty_runs, int)
+ or not 0 <= empty_runs <= 20):
+ raise ValueError(
+ "adaptive_cadence.empty_runs_before_idle должен быть от 0 до 20")
+ if not isinstance(event_wake, bool):
+ raise ValueError("adaptive_cadence.event_wake должен быть true или false")
+ if busy is None and (threshold or empty_runs):
+ raise ValueError(
+ "adaptive_cadence.backlog_above/empty_runs_before_idle требуют "
+ "busy_recurrence")
+ return {
+ "idle_recurrence": idle,
+ "busy_recurrence": busy,
+ "backlog_above": threshold,
+ "empty_runs_before_idle": empty_runs,
+ "event_wake": event_wake,
+ }
+
+
+def _adaptive_cadence_status(queue: dict, matching: dict | None,
+ backlog: int | None) -> dict | None:
+ policy = _adaptive_cadence_policy(queue)
+ if policy is None:
+ return None
+ busy = bool(
+ policy["busy_recurrence"] is not None
+ and isinstance(backlog, int)
+ and backlog > policy["backlog_above"]
+ )
+ temporary = matching.get("temporary_recurrence") if matching else None
+ empty_count = int(matching.get("temporary_empty_count") or 0) if matching else 0
+ draining = bool(
+ not busy and policy["empty_runs_before_idle"]
+ and temporary == policy["busy_recurrence"]
+ )
+ mode = "busy" if busy else "draining" if draining else (
+ "event" if policy["event_wake"] else "idle")
+ return {
+ **policy,
+ "mode": mode,
+ "effective_recurrence": (
+ matching.get("effective_recurrence") if matching else None),
+ "empty_runs": empty_count,
+ "series_present": matching is not None,
+ }
+
+
+def _reconcile_adaptive_cadence(
+ queue: dict, matching: dict | None, backlog: int | None,
+ publication_guard: dict | None = None) -> dict | None:
+ """Apply cadence only during a successful live queue observation."""
+ policy = _adaptive_cadence_policy(queue)
+ if policy is None or matching is None or not isinstance(backlog, int):
+ return _adaptive_cadence_status(queue, matching, backlog)
+ boost = bool(
+ policy["busy_recurrence"] is not None
+ and backlog > policy["backlog_above"])
+ result = db.apply_pipeline_series_cadence(
+ int(matching["id"]),
+ idle_recurrence=policy["idle_recurrence"],
+ busy_recurrence=policy["busy_recurrence"],
+ boost=boost,
+ empty_runs_before_idle=policy["empty_runs_before_idle"],
+ publication_guard=publication_guard,
+ )
+ if result is not None:
+ matching.update({
+ "recurrence": result["base_recurrence"],
+ "effective_recurrence": result["effective_recurrence"],
+ "temporary_recurrence": result["temporary_recurrence"],
+ "temporary_empty_limit": result["temporary_empty_limit"],
+ "temporary_empty_count": result["temporary_empty_count"],
+ })
+ return _adaptive_cadence_status(queue, matching, backlog)
+
+
def _recommendation(item: dict, backlog: int, capacity: int,
current_interval: str | None, target_hours: float,
avg_duration_seconds: int | None = None) -> dict:
@@ -1461,18 +1580,52 @@ def _window_metrics(snapshots: list[dict], current: dict, series_ids: list[int],
churn_items = sum(1 for seq in sequences.values() if len(seq) >= 3)
queue_deltas = {}
+ queue_throughput = {}
for queue_id, queue in current.get("queues", {}).items():
- old = baseline.get("queues", {}).get(queue_id, {}).get("backlog", 0)
+ old_queue = baseline.get("queues", {}).get(queue_id, {})
+ old = old_queue.get("backlog", 0)
queue_deltas[queue_id] = int(queue.get("backlog", 0)) - int(old)
+ old_keys = {item.get("key") for item in old_queue.get("items", [])
+ if isinstance(item, dict) and item.get("key")}
+ new_keys = {item.get("key") for item in queue.get("items", [])
+ if isinstance(item, dict) and item.get("key")}
+ queue_throughput[queue_id] = (
+ len(old_keys - new_keys)
+ if (bool(old_queue.get("membership_complete"))
+ and bool(queue.get("membership_complete")))
+ else None
+ )
current_total = sum(q.get("backlog", 0) for q in current.get("queues", {}).values())
baseline_total = sum(q.get("backlog", 0) for q in baseline.get("queues", {}).values())
+ # Rates describe the observations we actually have. If the closest
+ # baseline is ten hours old, dividing its delta by a five-hour display
+ # window would overstate throughput by 2x.
+ measured_hours = coverage
+ complete_membership = all(
+ bool(queue.get("membership_complete"))
+ and bool(baseline.get("queues", {}).get(queue_id, {}).get(
+ "membership_complete"))
+ for queue_id, queue in current.get("queues", {}).items()
+ )
+
+ def hourly(value: int) -> float | None:
+ return round(value / measured_hours, 2) if measured_hours > 0 else None
return {
"hours": hours, "coverage_hours": round(min(coverage, hours), 1),
"complete": complete, "backlog_delta": current_total - baseline_total,
+ "backlog_delta_per_hour": hourly(current_total - baseline_total),
"entered": len(entered), "exited": len(exited), "moved": len(moved),
+ "entered_per_hour": hourly(len(entered)) if complete_membership else None,
+ "exited_per_hour": hourly(len(exited)) if complete_membership else None,
"transitions": transitions, "churn_items": churn_items,
"queue_deltas": queue_deltas,
+ "queue_deltas_per_hour": {
+ queue_id: hourly(delta) for queue_id, delta in queue_deltas.items()},
+ "queue_throughput": queue_throughput,
+ "queue_throughput_per_hour": {
+ queue_id: hourly(count) if count is not None else None
+ for queue_id, count in queue_throughput.items()},
"runs": db.pipeline_run_metrics(series_ids, now - timedelta(hours=hours)),
}
@@ -1561,6 +1714,113 @@ def _matching_queue(task) -> tuple[str, dict, dict] | None:
return None
+def worker_lane_policy() -> dict | None:
+ """Load generic, profile-owned worker lanes.
+
+ A lane is one concurrency slot with an ordered list of preferred queues.
+ ``borrow`` lets its idle slot execute work from another lane. Profiles
+ without this opt-in retain the historical global priority/FIFO scheduler.
+ """
+ profiles = _profiles()
+ lanes = []
+ seen = set()
+ for profile_id, profile in profiles.items():
+ if not isinstance(profile, dict):
+ raise ValueError(f"pipeline profile {profile_id} должен быть JSON-объектом")
+ scheduler = profile.get("scheduler")
+ if scheduler is None:
+ continue
+ if not isinstance(scheduler, dict) or set(scheduler) != {"lanes"}:
+ raise ValueError("scheduler должен содержать только массив lanes")
+ configured = scheduler.get("lanes")
+ if not isinstance(configured, list) or not configured:
+ raise ValueError("scheduler.lanes должен быть непустым массивом")
+ queue_ids = {
+ str(queue.get("id")) for queue in profile.get("queues", [])
+ if isinstance(queue, dict) and queue.get("id") is not None
+ }
+ for raw in configured:
+ if not isinstance(raw, dict) or set(raw) - {"id", "queues", "borrow"}:
+ raise ValueError(
+ "каждая scheduler lane допускает только id, queues и borrow")
+ lane_id = raw.get("id")
+ queues = raw.get("queues")
+ borrow = raw.get("borrow", True)
+ if (not isinstance(lane_id, str) or not lane_id.strip()
+ or not re.fullmatch(r"[A-Za-z0-9._-]+", lane_id)):
+ raise ValueError("scheduler lane id должен быть непустым safe-id")
+ qualified = f"{profile_id}:{lane_id}"
+ if qualified in seen:
+ raise ValueError(f"повтор scheduler lane id: {qualified}")
+ seen.add(qualified)
+ if (not isinstance(queues, list) or not queues
+ or any(not isinstance(value, str) or not value.strip()
+ for value in queues)
+ or len(set(queues)) != len(queues)):
+ raise ValueError(
+ f"scheduler lane {qualified}.queues должен быть "
+ "непустым массивом уникальных id")
+ unknown = sorted(set(queues) - queue_ids)
+ if unknown:
+ raise ValueError(
+ f"scheduler lane {qualified} ссылается на неизвестные "
+ f"очереди: {', '.join(unknown)}")
+ if not isinstance(borrow, bool):
+ raise ValueError(f"scheduler lane {qualified}.borrow должен быть boolean")
+ lanes.append({
+ "id": qualified, "profile_id": profile_id,
+ "queues": list(queues), "borrow": borrow,
+ "order": len(lanes),
+ })
+ return {"profiles": profiles, "lanes": lanes} if lanes else None
+
+
+def worker_lane_rank(task, policy: dict, busy_lane_ids=()) -> tuple[tuple, str] | None:
+ """Return the best free lane and deterministic rank for one task.
+
+ Preferred work always beats borrowed work in a free slot. Within a lane,
+ the configured queue order is authoritative; ordinary task priority and
+ FIFO order break ties. ``None`` means no configured slot can claim now.
+ """
+ busy = set(busy_lane_ids or ())
+ free = [lane for lane in policy.get("lanes", []) if lane["id"] not in busy]
+ if not free:
+ return None
+ title = (getattr(task, "series_title", None)
+ or str(getattr(task, "prompt", "")).splitlines()[0]).lower()
+ matched = None
+ if getattr(task, "series_id", None):
+ for profile_id, profile in policy.get("profiles", {}).items():
+ for queue in profile.get("queues", []):
+ marker = str(queue.get("series_contains") or "").lower()
+ if marker and marker in title:
+ matched = (profile_id, str(queue.get("id")))
+ break
+ if matched:
+ break
+ priority = int(getattr(task, "priority", 5))
+ created = str(getattr(task, "created_at", ""))
+ task_id = int(getattr(task, "id", 0))
+ preferred = []
+ if matched:
+ profile_id, queue_id = matched
+ for lane in free:
+ if lane["profile_id"] != profile_id or queue_id not in lane["queues"]:
+ continue
+ preferred.append((
+ (0, lane["queues"].index(queue_id), lane["order"],
+ priority, created, task_id),
+ lane["id"],
+ ))
+ if preferred:
+ return min(preferred, key=lambda item: item[0])
+ borrowing = [lane for lane in free if lane["borrow"]]
+ if not borrowing:
+ return None
+ lane = min(borrowing, key=lambda item: item["order"])
+ return ((1, priority, created, task_id, lane["order"]), lane["id"])
+
+
def _diagnostic_match_count(diagnostics: dict, condition: dict) -> int:
field = condition.get("field")
if not isinstance(field, str):
@@ -2073,9 +2333,16 @@ def execution_route(task, fallback_prompt: str, working_dir: str | None = None,
preflight.get("reason") or preflight.get("error") or preflight_action
)
if preflight_action in {"empty", "wait"}:
+ preflight_verdict = str(preflight.get("verdict") or "ПУСТО").strip() or "ПУСТО"
+ # Only a validated targeted fallback route can authorize the silent,
+ # immediate stale outcome. An ordinary project tool returning
+ # empty/wait must never mint that capability through an arbitrary JSON
+ # field.
+ if preflight_verdict.upper() == "УСТАРЕЛО":
+ preflight_verdict = "НЕ СМОГ"
return {
"action": "complete_empty", "mode": "tool", "reason": preflight_reason,
- "verdict": str(preflight.get("verdict") or "ПУСТО"),
+ "verdict": preflight_verdict,
"profile_id": profile_id, "queue_id": queue.get("id"),
"preflight": preflight,
}
@@ -2121,8 +2388,12 @@ def execution_route(task, fallback_prompt: str, working_dir: str | None = None,
"разбери и выведи JSON; запрещено бросать исключение только по exit code "
"до разбора ответа. При ненулевом коде, неверном JSON, action не равном "
"validated или несовпадении цели остановись, не повторяй gate_command и "
- "next, и закончи строкой ИТОГ: НЕ СМОГ (gate-fallback: <точный error, "
- "reason или сырой ответ>). "
+ "next. Если структурированный ответ доказывает только смену exact target, "
+ "HEAD, executable-позиции или истечение lease до первой мутации, закончи "
+ "строкой ИТОГ: УСТАРЕЛО (gate-fallback: <точный error/reason>): "
+ "PromptPilot немедленно и безопасно перевыберет цель. При любой другой "
+ "причине закончи ИТОГ: НЕ СМОГ (gate-fallback: <точный error, reason "
+ "или сырой ответ>). "
"Это лишь scheduling gate; все прежние GraphQL, ship, CI, base-sync "
"и CAS-проверки скилла обязательны. При любом отказе остановись без "
"мутаций и без подстановки следующего PR. Один envelope — один PR.\n\n"
@@ -2449,6 +2720,8 @@ def _refresh_local_state(result: dict, profile: dict, series: list[dict], *,
}
queue["wake"] = (_wake_status(data["profile_id"], config, diagnostics)
if isinstance(diagnostics, dict) else None)
+ queue["adaptive_cadence"] = _adaptive_cadence_status(
+ config, matching, backlog if isinstance(backlog, int) else None)
bottleneck = max(
(queue for queue in data.get("queues", [])
@@ -2487,7 +2760,11 @@ def _refresh_local_state(result: dict, profile: dict, series: list[dict], *,
window = {
"hours": hours, "coverage_hours": 0, "complete": False,
"backlog_delta": 0, "entered": 0, "exited": 0, "moved": 0,
+ "backlog_delta_per_hour": None,
+ "entered_per_hour": None, "exited_per_hour": None,
"transitions": 0, "churn_items": 0, "queue_deltas": {},
+ "queue_deltas_per_hour": {}, "queue_throughput": {},
+ "queue_throughput_per_hour": {},
}
history[key] = window
if hours == 5:
@@ -2495,6 +2772,20 @@ def _refresh_local_state(result: dict, profile: dict, series: list[dict], *,
elif not isinstance(window.get("runs"), dict):
window["runs"] = dict(empty_runs)
+ safe_deferrals = sum(
+ 1 for item in matching_series
+ if item.get("next_status") in {"pending", "rate_limited"}
+ and bool(item.get("next_error"))
+ )
+ data["outcomes"] = {
+ "safe_deferrals_now": safe_deferrals,
+ "stale_reselections_5h": aggregate_runs.get("stale", 0),
+ "real_errors_5h": (
+ aggregate_runs.get("unresolved_unable", aggregate_runs.get("unable", 0))
+ + aggregate_runs.get("unresolved_failed", aggregate_runs.get("failed", 0))
+ ),
+ }
+
runtime = _pipeline_runtime(matching_series, now)
backlog_total = data.get("backlog_total")
health = _health(
@@ -2923,6 +3214,7 @@ def _analyze_without_budget(profile_id: str, series: list[dict], *,
paused_series = 0
matching_series = []
all_items = {}
+ cadence_reconciliations = []
for item in profile["queues"]:
queries = item.get("queries") or [item["query"]]
@@ -2965,6 +3257,8 @@ def _analyze_without_budget(profile_id: str, series: list[dict], *,
series_ids.append(matching["id"])
broken_series += int(bool(matching.get("broken")))
paused_series += int(bool(matching.get("paused")))
+ cadence = _adaptive_cadence_status(item, matching, backlog)
+ cadence_reconciliations.append((item, matching, backlog))
runs_needed = round(backlog / capacity, 1)
recommendation = _recommendation(
item, backlog, capacity,
@@ -2996,6 +3290,7 @@ def _analyze_without_budget(profile_id: str, series: list[dict], *,
if priority_settings else [],
"execution": _execution_status(item, matching.get("working_dir") if matching else None),
"wake": _wake_status(profile_id, item, diagnostics),
+ "adaptive_cadence": cadence,
**recommendation,
})
snapshot_queues[item["id"]] = {
@@ -3077,6 +3372,32 @@ def _analyze_without_budget(profile_id: str, series: list[dict], *,
raise _GitHubScanLeaseLost(
"SQLite lease/snapshot fence rejected GitHub scan publication")
db.prune_pipeline_snapshots(now - timedelta(days=31))
+ # Cadence changes are local mutations derived from this observation.
+ # Apply them only after both the full-cache CAS and the fenced history
+ # snapshot accepted the observation; a stale/losing scan must never
+ # retime a live series.
+ publication_guard = {
+ "epoch_key": _CACHE_EPOCH_KEY,
+ "epoch": str(cache_epoch),
+ "epoch_default": "0",
+ "profile_key": _published_profile_key(profile_id),
+ "profile_hash": profile_hash,
+ "revision_key": _published_profile_revision_key(profile_id),
+ "revision": str(refresh_revision),
+ }
+ try:
+ current_profile = _profiles().get(profile_id)
+ profile_still_current = (
+ isinstance(current_profile, dict)
+ and _profile_fingerprint(current_profile) == profile_hash
+ )
+ except (AttributeError, OSError, TypeError, ValueError, json.JSONDecodeError):
+ profile_still_current = False
+ if profile_still_current:
+ for item, matching, backlog in cadence_reconciliations:
+ _reconcile_adaptive_cadence(
+ item, matching, backlog,
+ publication_guard=publication_guard)
return _refresh_local_state(
result, profile, series, source="live",
generated_at=float(result["generated_at"]), entry_epoch=cache_epoch,
diff --git a/promptpilot/static/index.html b/promptpilot/static/index.html
index 853b33c..56276a7 100644
--- a/promptpilot/static/index.html
+++ b/promptpilot/static/index.html
@@ -2650,10 +2650,12 @@
План этапов
const h7d = data.history['168h'] || {complete:false, coverage_hours:0};
const h30d = data.history['720h'] || {complete:false, coverage_hours:0};
const signed = n => n > 0 ? `+${n}` : String(n);
+ const perHour = n => n == null ? '—' : `${signed(Number(n).toFixed(2))}/ч`;
const age = value => value == null ? '—' : value < 24 ? `${Math.round(value)} ч` : `${(value/24).toFixed(1)} д`;
const duration = seconds => seconds == null ? '—' : seconds < 3600 ? `${Math.round(seconds/60)} мин` : `${(seconds/3600).toFixed(1)} ч`;
const historyNote = h5.complete ? 'за последние 5 часов' : `история ${h5.coverage_hours} из 5 ч`;
const runs = h5.runs;
+ const outcomes = data.outcomes || {safe_deferrals_now:0, stale_reselections_5h:0, real_errors_5h:0};
const activeUnable = runs.unresolved_unable ?? runs.unable;
const activeFailed = runs.unresolved_failed ?? runs.failed;
const recoveredRuns = (runs.recovered_unable || 0) + (runs.recovered_failed || 0);
@@ -2727,7 +2729,9 @@ План этапов
const coreLimit = githubLimits.core;
const coreShare = coreLimit && coreLimit.limit ? coreLimit.remaining / coreLimit.limit : null;
const coreColor = coreShare == null ? 'var(--muted)' : coreShare <= 0.05 ? 'var(--red)' : coreShare <= 0.2 ? 'var(--yellow)' : '';
- const deltaByQueue = h5.queue_deltas || {};
+ const deltaRateByQueue = h5.queue_deltas_per_hour || {};
+ const throughputRateByQueue = h5.queue_throughput_per_hour || {};
+ const bottleneckQueue = (data.queues || []).find(q => q.id === data.bottleneck);
const diagnostics = data.diagnostics;
const runtime = data.runtime || {};
const healthLabel = data.health.label === 'нужна проверка' ? 'есть ожидания' : data.health.label;
@@ -2746,16 +2750,20 @@ План этапов
const lastLine = last && last.at
? `последний #${last.task_id}: ${timeAgo(last.at)} · ${last.verdict || last.status}${lastTokens}`
: 'завершённых запусков нет';
- const stageLine = `5 ч: ${stageRuns.runs || 0} запусков · действий ${stageRuns.ready || 0} · холостых ${stageRuns.no_change || 0}`;
+ const stageLine = `5 ч: ${stageRuns.runs || 0} запусков · действий ${stageRuns.ready || 0} · холостых ${stageRuns.no_change || 0} · фактический выход ${perHour(throughputRateByQueue[q.id])}`;
const route = q.execution || {configured:'skill',effective:'skill'};
const routeLine = `маршрут: ${route.configured}${route.configured!==route.effective ? ` → ${route.effective}` : ''}${route.reason && route.effective!=='tool' ? ` · ${route.reason}` : ''}`;
const wake = q.wake;
const wakeLine = wake ? (wake.suppressed
? `автозапуск: ждёт изменения состояния · снимок ${wake.fingerprint}`
: wake.ready ? `автозапуск: готов · ${wake.matches}` : 'автозапуск: работы нет') : null;
+ const cadence = q.adaptive_cadence;
+ const cadenceLine = cadence
+ ? `интервал: ${cadence.mode} · ${cadence.effective_recurrence || 'нет серии'}${cadence.empty_runs_before_idle ? ` · ПУСТО ${cadence.empty_runs}/${cadence.empty_runs_before_idle}` : ''}${cadence.event_wake ? ' · event wake' : ''}`
+ : null;
const bottleneck = data.bottleneck===q.id ? 'bottleneck' : '';
- return `| ${data.bottleneck===q.id?'⚠ ':''}${esc(q.title)} | ${q.backlog == null ? '—' : q.backlog} | ${h5.complete?signed(deltaByQueue[q.id]||0):'—'} | ${q.capacity} | ${q.runs_needed == null ? '—' : q.runs_needed} | ${duration(q.avg_duration_seconds)} | ${q.eta_hours == null ? '—' : q.eta_hours+' ч'} | ${age(q.age.oldest_hours)} | ${esc(q.interval || 'нет серии')} | ${esc(q.recommendation || '—')} |
- ${esc(lastLine)}${esc(stageLine)}${esc(routeLine)}${wakeLine?`${esc(wakeLine)}`:''} |
`;
+ return `| ${data.bottleneck===q.id?'⚠ ':''}${esc(q.title)} | ${q.backlog == null ? '—' : q.backlog} | ${h5.complete?perHour(deltaRateByQueue[q.id] ?? 0):'—'} | ${q.capacity} | ${q.runs_needed == null ? '—' : q.runs_needed} | ${duration(q.avg_duration_seconds)} | ${q.eta_hours == null ? '—' : q.eta_hours+' ч'} | ${age(q.age.oldest_hours)} | ${esc(q.interval || 'нет серии')} | ${esc(q.recommendation || '—')} |
+ ${esc(lastLine)}${esc(stageLine)}${esc(routeLine)}${wakeLine?`${esc(wakeLine)}`:''}${cadenceLine?`${esc(cadenceLine)}`:''} |
`;
}).join('');
const priorityHtml = data.priority_control ? `
Приоритет элементов очереди · P0 выполняется первым, старые задачи повышаются каждые ${data.priority_control.aging_hours} ч до P1
@@ -2788,11 +2796,13 @@ План этапов
${cacheNotice}
Снимок GitHub
${esc(cacheValue)}
${esc(cacheNote)}
-
Backlog сейчас
${data.backlog_total == null ? '—' : data.backlog_total}
Δ 5 ч: ${h5.complete?signed(h5.backlog_delta):'копится история'}
-
Вход / выход
${h5.complete?`${h5.entered} / ${h5.exited}`:'—'}
${historyNote}
+
Backlog сейчас
${data.backlog_total == null ? '—' : data.backlog_total}
Δ/ч: ${h5.complete?perHour(h5.backlog_delta_per_hour):'копится история'}
+
Вход / выход
${h5.complete?`${perHour(h5.entered_per_hour)} / ${perHour(h5.exited_per_hour)}`:'—'}
фактическая скорость · ${historyNote}
+
Узкое место
${bottleneckQueue?esc(bottleneckQueue.title):'нет'}
${bottleneckQueue&&bottleneckQueue.eta_hours!=null?`ETA ${bottleneckQueue.eta_hours} ч · очередь ${bottleneckQueue.backlog}`:'активной очереди нет'}
Переходы / churn
${h5.complete?`${h5.transitions} / ${h5.churn_items}`:'—'}
между очередями / 3+ состояний
Возраст p90 / max
${age(data.age.p90_hours)} / ${age(data.age.oldest_hours)}
текущие элементы
-
Прогоны 5 ч
${runs.runs}
с действием ${runs.ready} · без действий ${runs.no_change || 0} · пусто ${runs.empty}
+
Прогоны 5 ч
${runs.runs}
с действием ${runs.ready} · без действий ${runs.no_change || 0} · пусто ${runs.empty} · перевыбор ${runs.stale || 0}
+
Безопасные ожидания / ошибки
${outcomes.safe_deferrals_now + outcomes.stale_reselections_5h} / ${outcomes.real_errors_5h}
ожидают сейчас ${outcomes.safe_deferrals_now} · stale-перевыборов 5 ч ${outcomes.stale_reselections_5h} · активных ошибок ${outcomes.real_errors_5h}
Требуют внимания
${runs.human + activeUnable + activeFailed}
человек ${runs.human} · не смог ${activeUnable} · упало ${activeFailed}${recoveredRuns ? ` · восстановлено ${recoveredRuns}` : ''}
Последнее REVIEW
${esc(lastReviewValue)}
${esc(lastReviewNote)}
Токены 5 ч
${esc(tokenValue)}
измерено запусков: ${runs.tokens_known_runs || 0} из ${runs.runs}
@@ -2805,7 +2815,7 @@
План этапов
${budgetMetric}
${diagnosticHtml}
- | Этап | Очередь | Δ 5 ч | За прогон | Прогонов | Средний запуск | ETA | Самый старый | Интервал | Рекомендация |
${rows}
+ | Этап | Очередь | Δ/ч | За прогон | Прогонов | Средний запуск | ETA | Самый старый | Интервал | Рекомендация |
${rows}
${priorityHtml}
`;
} catch(e) { box.innerHTML = `Анализ очереди недоступен: ${esc(e.message)}
`; }
diff --git a/promptpilot/worker.py b/promptpilot/worker.py
index 82eb83f..41a8b10 100644
--- a/promptpilot/worker.py
+++ b/promptpilot/worker.py
@@ -135,7 +135,9 @@ def is_rate_limited(text: str, exit_code: int) -> bool:
# The agent is asked to end with this line so a finished task says WHAT
# happened, not just that the process exited 0. Parsed whether or not we asked.
-VERDICTS = ("ГОТОВО", "УЖЕ СДЕЛАНО", "НУЖЕН ЧЕЛОВЕК", "НЕ СМОГ", "ПУСТО")
+VERDICTS = (
+ "ГОТОВО", "УЖЕ СДЕЛАНО", "НУЖЕН ЧЕЛОВЕК", "НЕ СМОГ", "ПУСТО",
+)
VERDICT_RE = re.compile(r"^[ \t>*#-]*ИТОГ:\s*(" + "|".join(VERDICTS) + r")\b", re.M | re.I)
VERDICT_INSTRUCTION = (
@@ -483,7 +485,12 @@ def _maybe_recur(task, failed: bool = False):
recurrence = series["effective_recurrence"] if series else task.recurrence
if task.series_id and series is None: # series was explicitly ended
return
- next_dt = db.parse_recurrence(recurrence)
+ # The durable series latch allows one immediate re-election after a
+ # validated targeted gate-fallback. Consecutive stale results use the
+ # ordinary cadence, so a moving target cannot create an unbounded hot loop.
+ next_dt = (datetime.now(timezone.utc)
+ if series and series.get("stale_reselect_immediate")
+ else db.parse_recurrence(recurrence))
if not next_dt:
return
from .models import TaskCreate
@@ -567,7 +574,8 @@ def _notify_requeued(task, next_run, reason: str):
def _execute_herdr_task(task, provider_cfg, host=None, machine=None, prompt_override=None,
- admission_complete=None, require_closing_verdict=False):
+ admission_complete=None, require_closing_verdict=False,
+ allow_targeted_stale=False):
"""Run the task in a live herdr session (providers with executor=herdr).
host is the ssh target of the machine the session lives on (None = local).
@@ -628,7 +636,8 @@ def on_started(_pane_id):
on_worktree=on_worktree, on_pane=on_pane,
on_started=on_started,
prompt_override=prompt_override,
- require_closing_verdict=require_closing_verdict)
+ require_closing_verdict=require_closing_verdict,
+ allow_targeted_stale=allow_targeted_stale)
if outcome.get("cancelled"):
db.clear_cancel_request(task.id)
@@ -665,6 +674,8 @@ def on_started(_pane_id):
return
verdict = outcome.get("verdict") or parse_verdict(outcome["output"])
+ if verdict == "УСТАРЕЛО" and not allow_targeted_stale:
+ verdict = "НЕ СМОГ"
if require_closing_verdict and not outcome.get("verdict"):
_retry_sqlite_busy(
lambda: db.mark_failed(
@@ -765,6 +776,7 @@ def _execute_task_body(task, admission_complete=None):
agent_prompt = effective_prompt(task)
require_closing_verdict = False
+ allow_targeted_stale = False
if task.series_id:
try:
from . import pipeline_insights
@@ -818,7 +830,12 @@ def _execute_task_body(task, admission_complete=None):
return
if route["action"] == "complete_empty":
reason = route["reason"]
- verdict = route.get("verdict") or "ПУСТО"
+ verdict = str(route.get("verdict") or "ПУСТО").strip() or "ПУСТО"
+ # Defence in depth: execution_route normally normalizes this, but
+ # persistence must not trust an injected/custom route. УСТАРЕЛО is
+ # reserved for the validated targeted fallback provider path.
+ if verdict.upper() == "УСТАРЕЛО":
+ verdict = "НЕ СМОГ"
_retry_sqlite_busy(
lambda: db.mark_completed(
task.id,
@@ -835,9 +852,18 @@ def _execute_task_body(task, admission_complete=None):
require_closing_verdict = bool(
route.get("profile_id") and route.get("queue_id")
)
+ gate_command = route.get("gate_command")
+ allow_targeted_stale = bool(
+ require_closing_verdict
+ and route.get("next_already_run") is True
+ and isinstance(gate_command, list)
+ and gate_command
+ and all(isinstance(value, str) and value for value in gate_command)
+ )
if require_closing_verdict:
from .herdr_exec import ensure_closing_verdict_contract
- agent_prompt = ensure_closing_verdict_contract(agent_prompt)
+ agent_prompt = ensure_closing_verdict_contract(
+ agent_prompt, allow_targeted_stale=allow_targeted_stale)
if route.get("fallback_reason"):
if route.get("next_already_run"):
print(" -> Pipeline tool selected target; continuing full skill: "
@@ -870,6 +896,7 @@ def _execute_task_body(task, admission_complete=None):
prompt_override=agent_prompt,
admission_complete=admission_complete,
require_closing_verdict=require_closing_verdict,
+ allow_targeted_stale=allow_targeted_stale,
)
finally:
# Startup failures have no on_started callback. Once their durable
@@ -1111,7 +1138,8 @@ def _execute_task_body(task, admission_complete=None):
if require_closing_verdict:
from .herdr_exec import _closing_workflow_verdict
- verdict = _closing_workflow_verdict(verdict_source)
+ verdict = _closing_workflow_verdict(
+ verdict_source, allow_targeted_stale=allow_targeted_stale)
if not verdict:
_retry_sqlite_busy(
lambda: db.mark_failed(
@@ -1398,6 +1426,41 @@ def begin(self):
return self._pending
+def _claim_next_task(busy_keys=(), busy_lane_ids=()):
+ """Atomically claim by configured lane preference, with safe fallback.
+
+ Invalid optional lane configuration must be visible but must not freeze the
+ legacy queue. The project preflight still owns every external mutation;
+ this function only chooses which already-runnable local task gets a slot.
+ """
+ from . import pipeline_insights
+
+ try:
+ policy = pipeline_insights.worker_lane_policy()
+ except (AttributeError, OSError, TypeError, ValueError,
+ json.JSONDecodeError) as exc:
+ print(f" !! pipeline lane scheduler unavailable: {exc}", flush=True)
+ policy = None
+ if policy is None:
+ return db.get_next_runnable(
+ busy_keys=busy_keys, key_fn=lock_key), None
+
+ assignments = {}
+
+ def rank(task):
+ candidate = pipeline_insights.worker_lane_rank(
+ task, policy, busy_lane_ids)
+ if candidate is None:
+ return None
+ score, lane_id = candidate
+ assignments[task.id] = lane_id
+ return score
+
+ task = db.get_next_runnable(
+ busy_keys=busy_keys, key_fn=lock_key, order_key_fn=rank)
+ return task, assignments.get(task.id) if task is not None else None
+
+
def _warm_pipeline_runtime():
"""Load pipeline routing code before a runnable task is claimed.
@@ -1479,6 +1542,7 @@ def publish_heartbeat():
pool = None
in_flight = {} # Future -> (lock key, exact claimed task attempt)
+ task_lanes = {} # task id -> configured scheduler lane
stuck_recoveries = {} # exact attempt -> retry state; keeps its lock key
admission_fence = _AdmissionFence()
short_on_memory = False
@@ -1487,7 +1551,11 @@ def publish_heartbeat():
pool = ThreadPoolExecutor(max_workers=CONCURRENCY, thread_name_prefix="pp-task")
def reap():
+ before = {task.id for _lock, task in in_flight.values()}
_reap_futures(in_flight, stuck_recoveries)
+ after = {task.id for _lock, task in in_flight.values()}
+ for task_id in before - after:
+ task_lanes.pop(task_id, None)
while running:
reap()
@@ -1526,7 +1594,10 @@ def reap():
busy_keys = [lk for lk, _task in in_flight.values() if lk]
busy_keys.extend(
item["lock"] for item in stuck_recoveries.values() if item["lock"])
- task = db.get_next_runnable(busy_keys=busy_keys, key_fn=lock_key)
+ task, lane_id = _claim_next_task(
+ busy_keys=busy_keys,
+ busy_lane_ids=task_lanes.values(),
+ )
if task is None:
time.sleep(POLL_INTERVAL)
continue
@@ -1544,6 +1615,9 @@ def reap():
_drain_stuck_recoveries(stuck_recoveries)
else:
admission_complete = admission_fence.begin()
+ if lane_id is not None:
+ task_lanes[task.id] = lane_id
+ print(f" -> Scheduler lane: {lane_id}", flush=True)
in_flight[pool.submit(
_execute_task_with_admission_fence,
task, admission_complete)] = (lock_key(task), task)
diff --git a/tests/test_herdr_workflow_completion.py b/tests/test_herdr_workflow_completion.py
index e17c3f1..1c1a0df 100644
--- a/tests/test_herdr_workflow_completion.py
+++ b/tests/test_herdr_workflow_completion.py
@@ -35,6 +35,22 @@ def test_closing_workflow_verdict_accepts_only_final_line():
" ИТОГ: ГОТОВО — задача выполнена, изменения и\n"
" проверки перечислены выше"
) == "ГОТОВО"
+
+
+def test_stale_reselection_requires_targeted_exact_gate_fallback_form():
+ assert _closing_workflow_verdict(
+ "Gate proved that the exact HEAD changed.\n"
+ "ИТОГ: УСТАРЕЛО (gate-fallback: PR HEAD changed before mutation)"
+ ) == "НЕ СМОГ"
+ assert _closing_workflow_verdict(
+ "Gate proved that the exact HEAD changed.\n"
+ "ИТОГ: УСТАРЕЛО (gate-fallback: PR HEAD changed before mutation)",
+ allow_targeted_stale=True,
+ ) == "УСТАРЕЛО"
+ assert _closing_workflow_verdict(
+ "Gate refused.\nИТОГ: УСТАРЕЛО (PR HEAD changed)",
+ allow_targeted_stale=True,
+ ) == "НЕ СМОГ"
assert _closing_workflow_verdict(
"Проверять нечего.\nИТОГ: ПУСТО (очередь пуста)"
) == "ПУСТО"
@@ -54,6 +70,15 @@ def test_closing_workflow_verdict_accepts_only_final_line():
) == "ГОТОВО"
+def test_global_contract_does_not_advertise_targeted_stale_outcome():
+ generic = ensure_closing_verdict_contract("do work")
+ targeted = ensure_closing_verdict_contract(
+ "do targeted work", allow_targeted_stale=True)
+
+ assert "ИТОГ: УСТАРЕЛО" not in generic
+ assert "ИТОГ: УСТАРЕЛО (gate-fallback: точная причина)" in targeted
+
+
def test_agy_background_task_indicator_blocks_idle_completion():
assert _has_running_background_task(
"● [16:02:06] python -m pytest -q running"
diff --git a/tests/test_schedule_series.py b/tests/test_schedule_series.py
index 7381693..dac4e07 100644
--- a/tests/test_schedule_series.py
+++ b/tests/test_schedule_series.py
@@ -70,6 +70,297 @@ def test_fresh_install_has_no_project_specific_pipeline_profiles(tmp_path, monke
assert pipeline_insights.list_profiles() == []
+def test_first_valid_stale_verdict_reselects_immediately(isolated_db):
+ task = isolated_db.create_task(TaskCreate(
+ prompt="ExampleProject - REVIEW", recurrence="4h"))
+ claimed = isolated_db.get_next_runnable()
+ isolated_db.mark_completed(
+ claimed.id, "ИТОГ: УСТАРЕЛО (PR HEAD changed)", verdict="УСТАРЕЛО")
+
+ before = datetime.now(timezone.utc)
+ worker._recur_after_run(claimed)
+ series = isolated_db.get_series(task.series_id)
+
+ assert series["next_task_id"] != task.id
+ assert datetime.fromisoformat(series["next_run_at"]) <= \
+ before + timedelta(seconds=2)
+
+
+def test_second_consecutive_stale_uses_normal_cadence_and_non_stale_resets(
+ isolated_db):
+ task = isolated_db.create_task(TaskCreate(
+ prompt="ExampleProject - REVIEW", recurrence="4h"))
+ first = isolated_db.get_next_runnable()
+ isolated_db.mark_completed(
+ first.id, "ИТОГ: УСТАРЕЛО (gate-fallback: first)",
+ verdict="УСТАРЕЛО")
+ worker._recur_after_run(first)
+
+ second = isolated_db.get_next_runnable()
+ assert second is not None
+ isolated_db.mark_completed(
+ second.id, "ИТОГ: УСТАРЕЛО (gate-fallback: second)",
+ verdict="УСТАРЕЛО")
+
+ before = datetime.now(timezone.utc)
+ worker._recur_after_run(second)
+ series = isolated_db.get_series(task.series_id)
+
+ assert datetime.fromisoformat(series["next_run_at"]) >= \
+ before + timedelta(hours=3, minutes=59)
+
+ reset = isolated_db.prepare_series_recurrence(task.series_id, "ГОТОВО")
+ after_reset = isolated_db.prepare_series_recurrence(
+ task.series_id, "УСТАРЕЛО")
+ assert reset["stale_reselect_immediate"] is False
+ assert after_reset["stale_reselect_immediate"] is True
+
+
+def test_adaptive_cadence_uses_busy_interval_then_two_empty_runs(
+ isolated_db):
+ task = isolated_db.create_task(TaskCreate(
+ prompt="ExampleProject - TRIAGE", recurrence="30m",
+ scheduled_at=datetime.now(timezone.utc) + timedelta(hours=2)))
+
+ active = isolated_db.apply_pipeline_series_cadence(
+ task.series_id, idle_recurrence="30m", busy_recurrence="15m",
+ boost=True, empty_runs_before_idle=2)
+ first_empty = isolated_db.prepare_series_recurrence(task.series_id, "ПУСТО")
+ second_empty = isolated_db.prepare_series_recurrence(task.series_id, "ПУСТО")
+
+ assert active["effective_recurrence"] == "15m"
+ assert first_empty["effective_recurrence"] == "15m"
+ assert first_empty["temporary_empty_count"] == 1
+ assert second_empty["effective_recurrence"] == "30m"
+ assert second_empty["temporary_recurrence"] is None
+
+
+def test_adaptive_fix_cadence_returns_to_idle_at_threshold(isolated_db):
+ task = isolated_db.create_task(TaskCreate(
+ prompt="ExampleProject - FIX", recurrence="30m"))
+ matching = isolated_db.get_series(task.series_id)
+ queue = {
+ "adaptive_cadence": {
+ "idle_recurrence": "30m", "busy_recurrence": "15m",
+ "backlog_above": 3,
+ },
+ }
+
+ busy = pipeline_insights._reconcile_adaptive_cadence(queue, matching, 4)
+ idle = pipeline_insights._reconcile_adaptive_cadence(queue, matching, 3)
+
+ assert busy["mode"] == "busy"
+ assert busy["effective_recurrence"] == "15m"
+ assert idle["mode"] == "idle"
+ assert idle["effective_recurrence"] == "30m"
+
+
+def test_busy_observation_resets_temporary_empty_streak(isolated_db):
+ task = isolated_db.create_task(TaskCreate(
+ prompt="ExampleProject - TRIAGE", recurrence="30m"))
+ isolated_db.apply_pipeline_series_cadence(
+ task.series_id, idle_recurrence="30m", busy_recurrence="15m",
+ boost=True, empty_runs_before_idle=2)
+ empty = isolated_db.prepare_series_recurrence(task.series_id, "ПУСТО")
+
+ busy = isolated_db.apply_pipeline_series_cadence(
+ task.series_id, idle_recurrence="30m", busy_recurrence="15m",
+ boost=True, empty_runs_before_idle=2)
+
+ assert empty["temporary_empty_count"] == 1
+ assert busy["temporary_empty_count"] == 0
+ assert busy["effective_recurrence"] == "15m"
+
+
+def test_event_only_policy_clears_previous_temporary_boost(isolated_db):
+ task = isolated_db.create_task(TaskCreate(
+ prompt="ExampleProject - PLAN", recurrence="4h"))
+ assert isolated_db.update_series(task.series_id, {
+ "temporary_recurrence": "10m", "temporary_empty_limit": 3,
+ })
+
+ result = isolated_db.apply_pipeline_series_cadence(
+ task.series_id, idle_recurrence="4h", busy_recurrence=None)
+
+ assert result["effective_recurrence"] == "4h"
+ assert result["temporary_recurrence"] is None
+ assert result["temporary_empty_count"] == 0
+
+
+def test_adaptive_cadence_rejects_stale_profile_revision_guard(isolated_db):
+ task = isolated_db.create_task(TaskCreate(
+ prompt="ExampleProject - FIX", recurrence="30m"))
+ profile_key = "test:published-profile"
+ revision_key = "test:published-revision"
+ epoch_key = "test:cache-epoch"
+ isolated_db.set_setting(profile_key, "profile-a")
+ isolated_db.set_setting(revision_key, "7")
+ guard = {
+ "epoch_key": epoch_key, "epoch": "0", "epoch_default": "0",
+ "profile_key": profile_key, "profile_hash": "profile-a",
+ "revision_key": revision_key, "revision": "7",
+ }
+ accepted = isolated_db.apply_pipeline_series_cadence(
+ task.series_id, idle_recurrence="30m", busy_recurrence="15m",
+ boost=True, publication_guard=guard)
+ isolated_db.set_setting(revision_key, "8")
+
+ rejected = isolated_db.apply_pipeline_series_cadence(
+ task.series_id, idle_recurrence="30m", busy_recurrence="15m",
+ boost=False, publication_guard=guard)
+
+ assert accepted["effective_recurrence"] == "15m"
+ assert rejected is None
+ assert isolated_db.get_series(task.series_id)["effective_recurrence"] == "15m"
+
+
+def test_series_completion_and_live_cadence_write_are_serialized(
+ isolated_db, monkeypatch):
+ task = isolated_db.create_task(TaskCreate(
+ prompt="ExampleProject - TRIAGE", recurrence="30m"))
+ isolated_db.apply_pipeline_series_cadence(
+ task.series_id, idle_recurrence="30m", busy_recurrence="15m",
+ boost=True, empty_runs_before_idle=2)
+ real_connect = isolated_db._connect
+ prepare_locked = threading.Event()
+ release_prepare = threading.Event()
+ cadence_done = threading.Event()
+ errors = []
+
+ @contextmanager
+ def gated_connect(*args, **kwargs):
+ with real_connect(*args, **kwargs) as conn:
+ if threading.current_thread().name == "prepare-series":
+ assert kwargs.get("immediate") is True
+ prepare_locked.set()
+ if not release_prepare.wait(5):
+ raise TimeoutError("test did not release series transaction")
+ yield conn
+
+ monkeypatch.setattr(isolated_db, "_connect", gated_connect)
+
+ def prepare():
+ try:
+ isolated_db.prepare_series_recurrence(task.series_id, "ПУСТО")
+ except BaseException as exc:
+ errors.append(exc)
+
+ def observe_busy():
+ try:
+ isolated_db.apply_pipeline_series_cadence(
+ task.series_id, idle_recurrence="30m",
+ busy_recurrence="15m", boost=True,
+ empty_runs_before_idle=2)
+ except BaseException as exc:
+ errors.append(exc)
+ finally:
+ cadence_done.set()
+
+ preparing = threading.Thread(target=prepare, name="prepare-series")
+ observing = threading.Thread(target=observe_busy, name="observe-busy")
+ preparing.start()
+ assert prepare_locked.wait(5)
+ observing.start()
+ try:
+ assert cadence_done.wait(0.1) is False
+ finally:
+ release_prepare.set()
+ preparing.join(5)
+ observing.join(5)
+
+ assert not preparing.is_alive()
+ assert not observing.is_alive()
+ assert errors == []
+ assert isolated_db.get_series(task.series_id)["temporary_empty_count"] == 0
+
+
+def test_window_metrics_report_actual_hourly_rates(isolated_db, monkeypatch):
+ now = datetime.now(timezone.utc)
+ baseline = {
+ "queues": {"review": {
+ "backlog": 3,
+ "items": [{"key": "pr:1"}, {"key": "pr:2"}, {"key": "pr:3"}],
+ "membership_complete": True,
+ }},
+ }
+ current = {
+ "queues": {"review": {
+ "backlog": 2,
+ "items": [{"key": "pr:3"}, {"key": "pr:4"}],
+ "membership_complete": True,
+ }},
+ }
+ snapshots = [
+ {"captured_at": (now - timedelta(hours=5)).isoformat(),
+ "payload": baseline},
+ {"captured_at": now.isoformat(), "payload": current},
+ ]
+ monkeypatch.setattr(
+ isolated_db, "pipeline_run_metrics", lambda *_args, **_kwargs: {})
+
+ metrics = pipeline_insights._window_metrics(
+ snapshots, current, [], now, 5)
+
+ assert metrics["backlog_delta_per_hour"] == -0.2
+ assert metrics["entered_per_hour"] == 0.2
+ assert metrics["exited_per_hour"] == 0.4
+ assert metrics["queue_throughput_per_hour"]["review"] == 0.4
+
+
+def test_window_rates_use_actual_elapsed_baseline_gap(isolated_db, monkeypatch):
+ now = datetime.now(timezone.utc)
+ baseline = {"queues": {"review": {
+ "backlog": 3, "items": [{"key": "pr:1"}],
+ "membership_complete": True,
+ }}}
+ current = {"queues": {"review": {
+ "backlog": 2, "items": [], "membership_complete": True,
+ }}}
+ snapshots = [
+ {"captured_at": (now - timedelta(hours=10)).isoformat(),
+ "payload": baseline},
+ {"captured_at": now.isoformat(), "payload": current},
+ ]
+ monkeypatch.setattr(
+ isolated_db, "pipeline_run_metrics", lambda *_args, **_kwargs: {})
+
+ metrics = pipeline_insights._window_metrics(
+ snapshots, current, [], now, 5)
+
+ assert metrics["coverage_hours"] == 5
+ assert metrics["backlog_delta_per_hour"] == -0.1
+ assert metrics["queue_throughput_per_hour"]["review"] == 0.1
+
+
+def test_incomplete_membership_does_not_claim_set_based_throughput(
+ isolated_db, monkeypatch):
+ now = datetime.now(timezone.utc)
+ baseline = {"queues": {"review": {
+ "backlog": 3, "items": [{"key": "pr:1"}, {"key": "pr:2"}],
+ "membership_complete": False,
+ }}}
+ current = {"queues": {"review": {
+ "backlog": 2, "items": [{"key": "pr:2"}],
+ "membership_complete": True,
+ }}}
+ snapshots = [
+ {"captured_at": (now - timedelta(hours=5)).isoformat(),
+ "payload": baseline},
+ {"captured_at": now.isoformat(), "payload": current},
+ ]
+ monkeypatch.setattr(
+ isolated_db, "pipeline_run_metrics", lambda *_args, **_kwargs: {})
+
+ metrics = pipeline_insights._window_metrics(
+ snapshots, current, [], now, 5)
+
+ assert metrics["backlog_delta_per_hour"] == -0.2
+ assert metrics["entered_per_hour"] is None
+ assert metrics["exited_per_hour"] is None
+ assert metrics["queue_throughput"]["review"] is None
+ assert metrics["queue_throughput_per_hour"]["review"] is None
+
+
def test_project_health_check_is_immediate_and_accepts_red_json(monkeypatch):
calls = []
@@ -880,6 +1171,41 @@ def test_older_cross_process_refresh_cannot_overwrite_newer_cache(
assert restored["backlog_total"] == 2
+def test_older_different_profile_hash_cannot_reclaim_active_publication(
+ isolated_db):
+ old_profile = {
+ "title": "Example", "repository": "owner/example",
+ "queues": [{
+ "id": "review", "title": "Review", "capacity": 1,
+ "query": "label:old", "series_contains": "REVIEW",
+ }],
+ }
+ new_profile = {
+ **old_profile,
+ "queues": [{**old_profile["queues"][0], "query": "label:new"}],
+ }
+ pipeline_insights._cache.clear()
+ _cached, generation = pipeline_insights._cache_snapshot("profile-race")
+ epoch = pipeline_insights._cache_epoch()
+ newer = pipeline_insights._empty_cached_result(
+ "profile-race", new_profile)
+ newer["generated_at"] = 200.0
+ older = pipeline_insights._empty_cached_result(
+ "profile-race", old_profile)
+ older["generated_at"] = 100.0
+
+ assert pipeline_insights._publish_cache(
+ "profile-race", new_profile, generation, epoch, 2, newer) is True
+ assert pipeline_insights._publish_cache(
+ "profile-race", old_profile, generation, epoch, 1, older) is False
+ assert isolated_db.get_setting(
+ pipeline_insights._published_profile_key("profile-race")
+ ) == pipeline_insights._profile_fingerprint(new_profile)
+ assert isolated_db.get_setting(
+ pipeline_insights._published_profile_revision_key("profile-race")
+ ) == "2"
+
+
def test_stale_profile_analysis_uses_separate_durable_namespace(
isolated_db, monkeypatch):
old_profile = {
@@ -1445,11 +1771,19 @@ def test_invalidation_during_analysis_rejects_stale_full_cache_publish(isolated_
"queues": [{
"id": "review", "title": "Review", "capacity": 1,
"query": "is:pr", "series_contains": "REVIEW",
+ "adaptive_cadence": {
+ "idle_recurrence": "30m", "busy_recurrence": "15m",
+ "backlog_above": 0,
+ },
}],
}
+ isolated_db.create_task(TaskCreate(
+ prompt="Example - REVIEW", recurrence="30m"))
+ series = isolated_db.list_series()
entered_search = threading.Event()
release_search = threading.Event()
search_calls = []
+ cadence_calls = []
errors = []
def fake_search(repository, query):
@@ -1462,13 +1796,21 @@ def fake_search(repository, query):
def run_analysis():
try:
- pipeline_insights.analyze("cache-race", [], use_cache=False)
+ pipeline_insights.analyze("cache-race", series, use_cache=False)
except BaseException as exc: # preserve the worker-thread failure for the assertion
errors.append(exc)
+ real_reconcile = pipeline_insights._reconcile_adaptive_cadence
+
+ def track_reconcile(*args, **kwargs):
+ cadence_calls.append((args, kwargs))
+ return real_reconcile(*args, **kwargs)
+
monkeypatch.setattr(pipeline_insights, "_profiles", lambda: {"cache-race": profile})
monkeypatch.setattr(pipeline_insights, "_github_search", fake_search)
monkeypatch.setattr(pipeline_insights, "_run_profile_health_check", lambda _profile: None)
+ monkeypatch.setattr(
+ pipeline_insights, "_reconcile_adaptive_cadence", track_reconcile)
pipeline_insights._discard_cache()
analysis = threading.Thread(target=run_analysis)
analysis.start()
@@ -1483,15 +1825,64 @@ def run_analysis():
assert isolated_db.get_setting(
pipeline_insights._cache_key(
"cache-race", pipeline_insights._profile_fingerprint(profile))) is None
+ assert cadence_calls == []
- pipeline_insights.analyze("cache-race", [], use_cache=False)
+ pipeline_insights.analyze("cache-race", series, use_cache=False)
assert len(search_calls) == 2
+ assert len(cadence_calls) == 1
finally:
release_search.set()
analysis.join(5)
pipeline_insights._discard_cache()
+def test_invalidation_after_snapshot_fences_cadence_write(
+ isolated_db, monkeypatch):
+ profile = {
+ "title": "Example", "repository": "owner/example",
+ "queues": [{
+ "id": "fix", "title": "Fix", "capacity": 1,
+ "query": "is:issue", "series_contains": "FIX",
+ "adaptive_cadence": {
+ "idle_recurrence": "30m", "busy_recurrence": "15m",
+ "backlog_above": 0,
+ },
+ }],
+ }
+ task = isolated_db.create_task(TaskCreate(
+ prompt="Example - FIX", recurrence="30m"))
+ series = isolated_db.list_series()
+ real_prune = isolated_db.prune_pipeline_snapshots
+ invalidated = []
+
+ def invalidate_before_cadence(cutoff):
+ real_prune(cutoff)
+ pipeline_insights._discard_cache()
+ invalidated.append(True)
+
+ monkeypatch.setattr(
+ pipeline_insights, "_profiles", lambda: {"cadence-race": profile})
+ monkeypatch.setattr(pipeline_insights, "_github_search", lambda *_args: {
+ "count": 1, "items": [], "membership_complete": True,
+ })
+ monkeypatch.setattr(
+ pipeline_insights, "_run_profile_health_check", lambda _profile: None)
+ monkeypatch.setattr(
+ isolated_db, "prune_pipeline_snapshots", invalidate_before_cadence)
+ pipeline_insights._discard_cache()
+
+ try:
+ result = pipeline_insights.analyze(
+ "cadence-race", series, use_cache=False)
+
+ assert invalidated == [True]
+ assert result["cache"]["invalidated"] is True
+ assert isolated_db.get_series(task.series_id)[
+ "effective_recurrence"] == "30m"
+ finally:
+ pipeline_insights._discard_cache()
+
+
def test_cache_read_observes_invalidation_with_payload_at_one_snapshot(
isolated_db, monkeypatch):
profile = {
@@ -2460,6 +2851,39 @@ def test_pipeline_execution_empty_completes_without_provider(isolated_db, monkey
assert route["reason"] == "queue is empty"
+@pytest.mark.parametrize("action", ["empty", "wait"])
+def test_pipeline_execution_empty_cannot_authorize_stale(
+ isolated_db, monkeypatch, tmp_path, action):
+ helper = tmp_path / "pipelinectl.py"
+ helper.write_text(
+ "import json; print(json.dumps({"
+ f"'action':'{action}','verdict':'УСТАРЕЛО','reason':'generic tool'"
+ "}))",
+ encoding="utf-8",
+ )
+ task = isolated_db.create_task(TaskCreate(
+ prompt="ExampleProject - REVIEW\n/review-queue", recurrence="4h",
+ ))
+ profile = {
+ "title": "Example", "repository": "owner/example",
+ "queues": [{
+ "id": "review", "title": "Review", "query": "is:pr",
+ "series_contains": "ExampleProject - REVIEW",
+ "execution": {
+ "mode": "auto", "stage": "review",
+ "command": ["{python}", "pipelinectl.py", "next", "{stage}"],
+ "required_paths": ["pipelinectl.py"],
+ },
+ }],
+ }
+ monkeypatch.setattr(pipeline_insights, "_profiles", lambda: {"example": profile})
+
+ route = pipeline_insights.execution_route(task, task.prompt, str(tmp_path))
+
+ assert route["action"] == "complete_empty"
+ assert route["verdict"] == "НЕ СМОГ"
+
+
def test_worker_settles_preflight_empty_without_loading_provider(isolated_db, monkeypatch):
task = isolated_db.create_task(TaskCreate(
prompt="ExampleProject - REVIEW", recurrence="4h",
@@ -2486,6 +2910,33 @@ def test_worker_settles_preflight_empty_without_loading_provider(isolated_db, mo
assert "токены не потрачены" in settled.result
+def test_worker_never_persists_stale_from_complete_empty(isolated_db, monkeypatch):
+ task = isolated_db.create_task(TaskCreate(
+ prompt="ExampleProject - REVIEW", recurrence="4h",
+ ))
+ task = isolated_db.get_next_runnable()
+ monkeypatch.setattr(pipeline_insights, "dispatch_gate", lambda _task: None)
+ monkeypatch.setattr(
+ pipeline_insights, "execution_route",
+ lambda *_args, **_kwargs: {
+ "action": "complete_empty", "mode": "tool",
+ "reason": "generic tool", "verdict": "УСТАРЕЛО",
+ },
+ )
+ monkeypatch.setattr(
+ worker, "load_providers",
+ lambda: (_ for _ in ()).throw(AssertionError("provider must not be loaded")),
+ )
+
+ worker._execute_task_inner(task)
+
+ settled = isolated_db.get_task(task.id)
+ assert settled.status.value == "completed"
+ assert settled.verdict == "НЕ СМОГ"
+ assert "ИТОГ: НЕ СМОГ (generic tool)" in settled.result
+ assert "ИТОГ: УСТАРЕЛО" not in settled.result
+
+
def test_pipeline_execution_tool_fallback_skips_preflight_prompt(isolated_db, monkeypatch, tmp_path):
helper = tmp_path / "pipelinectl.py"
helper.write_text(
@@ -2788,6 +3239,21 @@ def test_pipeline_run_metrics_distinguish_semantic_failure_from_process_failure(
assert metrics["recovered_failed"] == 0
+def test_pipeline_run_metrics_count_stale_reselection_as_safe(isolated_db):
+ task = isolated_db.create_task(TaskCreate(
+ prompt="Project - REVIEW", recurrence="1h"))
+ isolated_db.mark_completed(
+ task.id, "ИТОГ: УСТАРЕЛО", verdict="УСТАРЕЛО")
+
+ metrics = isolated_db.pipeline_run_metrics(
+ [task.series_id], datetime.now(timezone.utc) - timedelta(hours=1))
+
+ assert metrics["runs"] == 1
+ assert metrics["stale"] == 1
+ assert metrics["unable"] == 0
+ assert metrics["failed"] == 0
+
+
def test_pipeline_run_metrics_clear_incident_after_later_success(isolated_db):
unable = isolated_db.create_task(TaskCreate(prompt="Project - REVIEW", recurrence="1h"))
ready = isolated_db.create_task(TaskCreate(
diff --git a/tests/test_worker_admission.py b/tests/test_worker_admission.py
index e490133..ee801c7 100644
--- a/tests/test_worker_admission.py
+++ b/tests/test_worker_admission.py
@@ -28,6 +28,93 @@ def test_admission_fence_blocks_next_claim_until_current_admission_finishes():
assert fence.wait(0) is False
+def test_worker_lane_policy_prefers_home_slots_then_borrows(monkeypatch):
+ from promptpilot import pipeline_insights
+
+ profile = {
+ "title": "Pipeline", "repository": "owner/repo",
+ "queues": [
+ {"id": "triage", "series_contains": " - TRIAGE"},
+ {"id": "plan", "series_contains": " - PLAN"},
+ {"id": "fix", "series_contains": " - FIX"},
+ {"id": "review", "series_contains": " - REVIEW"},
+ {"id": "merge", "series_contains": " - MERGE"},
+ ],
+ "scheduler": {"lanes": [
+ {"id": "integration", "queues": ["merge", "review"]},
+ {"id": "production", "queues": ["review", "fix"]},
+ {"id": "intake", "queues": ["triage", "plan"]},
+ ]},
+ }
+ monkeypatch.setattr(pipeline_insights, "_profiles", lambda: {"repo": profile})
+ policy = pipeline_insights.worker_lane_policy()
+
+ def task(number, stage):
+ return SimpleNamespace(
+ id=number, series_id=number, series_title=f"Repo - {stage}",
+ prompt=f"Repo - {stage}", priority=1,
+ created_at=datetime(2026, 1, number, tzinfo=timezone.utc),
+ )
+
+ assert pipeline_insights.worker_lane_rank(
+ task(1, "MERGE"), policy)[1] == "repo:integration"
+ assert pipeline_insights.worker_lane_rank(
+ task(2, "REVIEW"), policy)[1] == "repo:production"
+ assert pipeline_insights.worker_lane_rank(
+ task(3, "TRIAGE"), policy)[1] == "repo:intake"
+ assert pipeline_insights.worker_lane_rank(
+ task(4, "REVIEW"), policy, {"repo:production"})[1] == "repo:integration"
+ assert pipeline_insights.worker_lane_rank(
+ task(5, "REVIEW"), policy,
+ {"repo:integration", "repo:production"})[1] == "repo:intake"
+ assert pipeline_insights.worker_lane_rank(
+ task(6, "FIX"), policy,
+ {"repo:integration", "repo:production", "repo:intake"}) is None
+
+
+def test_lane_claim_is_atomic_and_queue_order_beats_fifo(
+ isolated_db, monkeypatch):
+ from promptpilot import pipeline_insights
+
+ profile = {
+ "title": "Pipeline", "repository": "owner/repo",
+ "queues": [
+ {"id": "fix", "series_contains": " - FIX"},
+ {"id": "review", "series_contains": " - REVIEW"},
+ {"id": "merge", "series_contains": " - MERGE"},
+ ],
+ "scheduler": {"lanes": [
+ {"id": "integration", "queues": ["merge", "review"]},
+ {"id": "production", "queues": ["review", "fix"]},
+ ]},
+ }
+ monkeypatch.setattr(pipeline_insights, "_profiles", lambda: {"repo": profile})
+ fix = isolated_db.create_task(TaskCreate(prompt="Repo - FIX", recurrence="1h"))
+ merge = isolated_db.create_task(TaskCreate(prompt="Repo - MERGE", recurrence="1h"))
+
+ claimed, lane = worker._claim_next_task()
+
+ assert claimed.id == merge.id
+ assert claimed.id != fix.id
+ assert lane == "repo:integration"
+
+
+def test_malformed_optional_profile_falls_back_to_legacy_fifo(
+ isolated_db, monkeypatch, capsys):
+ from promptpilot import pipeline_insights
+
+ first = isolated_db.create_task(TaskCreate(prompt="first", priority=3))
+ isolated_db.create_task(TaskCreate(prompt="second", priority=3))
+ monkeypatch.setattr(
+ pipeline_insights, "_profiles", lambda: {"broken": []})
+
+ claimed, lane = worker._claim_next_task()
+
+ assert claimed.id == first.id
+ assert lane is None
+ assert "pipeline lane scheduler unavailable" in capsys.readouterr().out
+
+
def test_worker_recovers_then_warms_pipeline_before_claiming(monkeypatch):
events = []
handlers = {}
@@ -101,6 +188,7 @@ def route(*_args, **_kwargs):
def provider(*_args, **kwargs):
assert admission_complete.is_set() is False
assert kwargs["require_closing_verdict"] is True
+ assert kwargs["allow_targeted_stale"] is False
calls.append("provider")
worker._signal_admission_complete(kwargs["admission_complete"])
assert admission_complete.is_set() is True
@@ -112,6 +200,59 @@ def provider(*_args, **kwargs):
assert calls == ["route", "provider"]
+def test_only_targeted_fallback_route_enables_stale_reselection(monkeypatch):
+ admission_complete = threading.Event()
+ task = SimpleNamespace(
+ id=42, series_id=8, prompt="Example - REVIEW",
+ provider="test-provider", working_dir=None, machine=None,
+ )
+ from promptpilot import pipeline_insights
+
+ monkeypatch.setattr(pipeline_insights, "dispatch_gate", lambda _task: None)
+ monkeypatch.setattr(
+ pipeline_insights, "execution_route",
+ lambda *_args, **_kwargs: {
+ "action": "prompt", "mode": "skill", "prompt": "run",
+ "profile_id": "example", "queue_id": "review",
+ "next_already_run": True,
+ "gate_command": ["python", "pipelinectl.py", "gate-fallback"],
+ },
+ )
+ monkeypatch.setattr(
+ worker, "load_providers",
+ lambda: {"test-provider": {"executor": "herdr"}},
+ )
+ captured = {}
+
+ def provider(*_args, **kwargs):
+ captured.update(kwargs)
+ worker._signal_admission_complete(kwargs["admission_complete"])
+
+ monkeypatch.setattr(worker, "_execute_herdr_task", provider)
+
+ worker._execute_task_body(task, admission_complete)
+
+ assert captured["require_closing_verdict"] is True
+ assert captured["allow_targeted_stale"] is True
+ assert "ИТОГ: УСТАРЕЛО (gate-fallback: точная причина)" in captured["prompt_override"]
+ assert worker.parse_verdict("ИТОГ: УСТАРЕЛО (gate-fallback: changed)") == ""
+
+
+def test_only_valid_routine_completion_verdicts_are_silent_in_telegram():
+ from promptpilot.bot import _silent_completion
+
+ completed = {"status": SimpleNamespace(value="completed")}
+
+ assert _silent_completion(SimpleNamespace(
+ **completed, verdict="ПУСТО")) is True
+ assert _silent_completion(SimpleNamespace(
+ **completed, verdict="УСТАРЕЛО")) is True
+ assert _silent_completion(SimpleNamespace(
+ **completed, verdict="НЕ СМОГ")) is False
+ assert _silent_completion(SimpleNamespace(
+ status=SimpleNamespace(value="failed"), verdict="УСТАРЕЛО")) is False
+
+
def test_sqlite_busy_control_write_is_retried(monkeypatch):
attempts = []
sleeps = []