From e09865addb965d9f9112b9c328dc09702434e36f Mon Sep 17 00:00:00 2001 From: ibrog Date: Thu, 17 Sep 2026 11:27:10 +0300 Subject: [PATCH 1/3] feat(pipeline): add adaptive throughput controls Generated-with: Codex --- README.md | 75 +++++++ promptpilot/bot.py | 29 ++- promptpilot/db.py | 128 +++++++++++- promptpilot/herdr_exec.py | 3 +- promptpilot/models.py | 2 +- promptpilot/pipeline_insights.py | 265 +++++++++++++++++++++++- promptpilot/static/index.html | 26 ++- promptpilot/worker.py | 68 +++++- tests/test_herdr_workflow_completion.py | 7 + tests/test_schedule_series.py | 142 ++++++++++++- tests/test_worker_admission.py | 71 +++++++ 11 files changed, 788 insertions(+), 28 deletions(-) diff --git a/README.md b/README.md index 568dd0e..36eaf9c 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,14 @@ pp note 42 --clear # убрать просыпающийся каждые два часа, превращает бота в будильник. Ошибки (`failed`) шлются всегда, независимо от итога. +Для распознанной серии конвейера контракт дополнительно разрешает `УСТАРЕЛО` — +отдельный безопасный исход: точная цель, HEAD или +её executable-позиция изменились **до первой мутации**. Это не `НЕ СМОГ` и не +ошибка здоровья. Следующее вхождение серии создаётся сразу, чтобы выполнить +свежий election; в дашборде такие перевыборы считаются отдельно. Неоднозначный +отказ, проблема доступа и ошибка после начатой мутации по-прежнему должны +заканчиваться `НЕ СМОГ` или `НУЖЕН ЧЕЛОВЕК`, а не быстрым циклом перевыбора. + Парсится он **всегда**, даже когда мы не просили — если агент сам закончил такой строкой, итог подхватится. Побеждает последнее совпадение: формат могли процитировать по дороге, а вердикт — это закрывающая строка. diff --git a/promptpilot/bot.py b/promptpilot/bot.py index b249d89..3e6062a 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) diff --git a/promptpilot/db.py b/promptpilot/db.py index ca9e622..eacac91 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,107 @@ 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) -> 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") + + 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 + 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: + reset_counter = ( + before.get("temporary_recurrence") != busy_recurrence + or before.get("temporary_until") is not None + or before.get("temporary_empty_limit") + != (empty_runs_before_idle or None) + ) + desired["temporary_recurrence"] = busy_recurrence + desired["temporary_until"] = None + desired["temporary_empty_limit"] = ( + empty_runs_before_idle or None) + if reset_counter: + 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 + + 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") @@ -1646,7 +1760,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 +1811,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 diff --git a/promptpilot/herdr_exec.py b/promptpilot/herdr_exec.py index b69780f..975c78f 100644 --- a/promptpilot/herdr_exec.py +++ b/promptpilot/herdr_exec.py @@ -70,10 +70,11 @@ ИТОГ: УЖЕ СДЕЛАНО (краткая причина) ИТОГ: НУЖЕН ЧЕЛОВЕК (краткая причина) ИТОГ: НЕ СМОГ (краткая причина) +ИТОГ: УСТАРЕЛО (цель или HEAD изменились до первой мутации; нужен перевыбор) ИТОГ: ПУСТО (краткая причина) {WORKFLOW_CONTRACT_END}""" WORKFLOW_CLOSING_VERDICT_RE = re.compile( - r"^ИТОГ:\s*(ГОТОВО|УЖЕ СДЕЛАНО|НУЖЕН ЧЕЛОВЕК|НЕ СМОГ|ПУСТО)" + r"^ИТОГ:\s*(ГОТОВО|УЖЕ СДЕЛАНО|НУЖЕН ЧЕЛОВЕК|НЕ СМОГ|УСТАРЕЛО|ПУСТО)" r"(?:\s*(?:[—-]\s*.*|\([^\r\n)]*\)))?$", re.IGNORECASE, ) 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..8227cc5 100644 --- a/promptpilot/pipeline_insights.py +++ b/promptpilot/pipeline_insights.py @@ -1355,6 +1355,107 @@ 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) -> 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"], + ) + 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 +1562,37 @@ 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) 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()) + measured_hours = min(coverage, hours) + + 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)), + "exited_per_hour": hourly(len(exited)), "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) for queue_id, count in queue_throughput.items()}, "runs": db.pipeline_run_metrics(series_ids, now - timedelta(hours=hours)), } @@ -1561,6 +1681,111 @@ 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(): + 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): @@ -2121,8 +2346,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 +2678,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 +2718,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 +2730,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 +3172,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 +3215,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 +3248,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 +3330,12 @@ 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. + for item, matching, backlog in cadence_reconciliations: + _reconcile_adaptive_cadence(item, matching, backlog) 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} - ${rows}
ЭтапОчередьΔ 5 чЗа прогонПрогоновСредний запускETAСамый старыйИнтервалРекомендация
+ ${rows}
ЭтапОчередьΔ/чЗа прогонПрогоновСредний запускETAСамый старыйИнтервалРекомендация
${priorityHtml} `; } catch(e) { box.innerHTML = `
Анализ очереди недоступен: ${esc(e.message)}
`; } diff --git a/promptpilot/worker.py b/promptpilot/worker.py index 82eb83f..92670ca 100644 --- a/promptpilot/worker.py +++ b/promptpilot/worker.py @@ -135,7 +135,10 @@ 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 +486,21 @@ 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) + # A stale exact target is neither a failed run nor useful work. Re-elect it + # immediately instead of making a busy queue wait for its ordinary cadence. + # Only the explicit closing verdict enables this path: generic failures and + # ambiguous gate errors keep their normal bounded schedule. + stale_pipeline_target = False + if (task.verdict or "").upper() == "УСТАРЕЛО": + try: + from . import pipeline_insights + stale_pipeline_target = pipeline_insights._matching_queue(task) is not None + except (AttributeError, OSError, TypeError, ValueError): + # A broken/temporarily unreadable optional profile must not turn + # an ordinary recurring task into an unbounded immediate loop. + stale_pipeline_target = False + next_dt = (datetime.now(timezone.utc) if stale_pipeline_target + else db.parse_recurrence(recurrence)) if not next_dt: return from .models import TaskCreate @@ -1398,6 +1415,40 @@ 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 (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 +1530,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 +1539,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 +1582,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 +1603,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..901661b 100644 --- a/tests/test_herdr_workflow_completion.py +++ b/tests/test_herdr_workflow_completion.py @@ -35,6 +35,13 @@ def test_closing_workflow_verdict_accepts_only_final_line(): " ИТОГ: ГОТОВО — задача выполнена, изменения и\n" " проверки перечислены выше" ) == "ГОТОВО" + + +def test_closing_workflow_verdict_accepts_stale_reselection_outcome(): + assert _closing_workflow_verdict( + "Gate proved that the exact HEAD changed.\n" + "ИТОГ: УСТАРЕЛО (PR HEAD changed before mutation)" + ) == "УСТАРЕЛО" assert _closing_workflow_verdict( "Проверять нечего.\nИТОГ: ПУСТО (очередь пуста)" ) == "ПУСТО" diff --git a/tests/test_schedule_series.py b/tests/test_schedule_series.py index 7381693..d992fb0 100644 --- a/tests/test_schedule_series.py +++ b/tests/test_schedule_series.py @@ -70,6 +70,111 @@ def test_fresh_install_has_no_project_specific_pipeline_profiles(tmp_path, monke assert pipeline_insights.list_profiles() == [] +def test_stale_pipeline_verdict_reselects_immediately(isolated_db, monkeypatch): + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": PIPELINE_PROFILE}) + 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_stale_verdict_keeps_normal_cadence_outside_pipeline( + isolated_db, monkeypatch): + monkeypatch.setattr(pipeline_insights, "_profiles", lambda: {}) + task = isolated_db.create_task(TaskCreate( + prompt="Ordinary recurring task", recurrence="4h")) + claimed = isolated_db.get_next_runnable() + isolated_db.mark_completed( + claimed.id, "ИТОГ: УСТАРЕЛО", verdict="УСТАРЕЛО") + + before = datetime.now(timezone.utc) + worker._recur_after_run(claimed) + series = isolated_db.get_series(task.series_id) + + assert datetime.fromisoformat(series["next_run_at"]) >= \ + before + timedelta(hours=3, minutes=59) + + +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_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"}], + }}, + } + current = { + "queues": {"review": { + "backlog": 2, + "items": [{"key": "pr:3"}, {"key": "pr:4"}], + }}, + } + 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_project_health_check_is_immediate_and_accepts_red_json(monkeypatch): calls = [] @@ -1445,11 +1550,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 +1575,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): + cadence_calls.append(args) + return real_reconcile(*args) + 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,9 +1604,11 @@ 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) @@ -2788,6 +2911,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..108c9c2 100644 --- a/tests/test_worker_admission.py +++ b/tests/test_worker_admission.py @@ -28,6 +28,77 @@ 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_worker_recovers_then_warms_pipeline_before_claiming(monkeypatch): events = [] handlers = {} From 6ab1bb735f4ee92b3276cc76e779f00f01501de3 Mon Sep 17 00:00:00 2001 From: ibrog Date: Thu, 17 Sep 2026 11:48:48 +0300 Subject: [PATCH 2/3] fix(pipeline): bound stale reselection and cadence writes Generated-with: Codex --- README.md | 16 +- promptpilot/bot.py | 17 +- promptpilot/db.py | 100 +++++++- promptpilot/herdr_exec.py | 46 +++- promptpilot/pipeline_insights.py | 75 +++++- promptpilot/worker.py | 49 ++-- tests/test_herdr_workflow_completion.py | 22 +- tests/test_schedule_series.py | 294 ++++++++++++++++++++++-- tests/test_worker_admission.py | 70 ++++++ 9 files changed, 611 insertions(+), 78 deletions(-) diff --git a/README.md b/README.md index 36eaf9c..effb227 100644 --- a/README.md +++ b/README.md @@ -1429,13 +1429,15 @@ pp note 42 --clear # убрать просыпающийся каждые два часа, превращает бота в будильник. Ошибки (`failed`) шлются всегда, независимо от итога. -Для распознанной серии конвейера контракт дополнительно разрешает `УСТАРЕЛО` — -отдельный безопасный исход: точная цель, HEAD или -её executable-позиция изменились **до первой мутации**. Это не `НЕ СМОГ` и не -ошибка здоровья. Следующее вхождение серии создаётся сразу, чтобы выполнить -свежий election; в дашборде такие перевыборы считаются отдельно. Неоднозначный -отказ, проблема доступа и ошибка после начатой мутации по-прежнему должны -заканчиваться `НЕ СМОГ` или `НУЖЕН ЧЕЛОВЕК`, а не быстрым циклом перевыбора. +Только 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 3e6062a..96e8bd7 100644 --- a/promptpilot/bot.py +++ b/promptpilot/bot.py @@ -3315,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. @@ -3360,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 eacac91..276c2c6 100644 --- a/promptpilot/db.py +++ b/promptpilot/db.py @@ -1092,7 +1092,8 @@ def update_series(series_id: int, fields: dict) -> bool: 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) -> Optional[dict]: + 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`` @@ -1114,7 +1115,33 @@ def apply_pipeline_series_cadence( 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() @@ -1130,23 +1157,25 @@ def apply_pipeline_series_cadence( and before.get("temporary_recurrence") == busy_recurrence ) if boost: - reset_counter = ( - before.get("temporary_recurrence") != busy_recurrence - or before.get("temporary_until") is not None - or before.get("temporary_empty_limit") - != (empty_runs_before_idle or None) - ) desired["temporary_recurrence"] = busy_recurrence desired["temporary_until"] = None desired["temporary_empty_limit"] = ( empty_runs_before_idle or None) - if reset_counter: - desired["temporary_empty_count"] = 0 + # 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 ( @@ -1555,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) @@ -1584,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 @@ -1925,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)): @@ -1946,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), @@ -1954,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 975c78f..9734331 100644 --- a/promptpilot/herdr_exec.py +++ b/promptpilot/herdr_exec.py @@ -70,14 +70,17 @@ ИТОГ: УЖЕ СДЕЛАНО (краткая причина) ИТОГ: НУЖЕН ЧЕЛОВЕК (краткая причина) ИТОГ: НЕ СМОГ (краткая причина) -ИТОГ: УСТАРЕЛО (цель или HEAD изменились до первой мутации; нужен перевыбор) ИТОГ: ПУСТО (краткая причина) {WORKFLOW_CONTRACT_END}""" WORKFLOW_CLOSING_VERDICT_RE = re.compile( - r"^ИТОГ:\s*(ГОТОВО|УЖЕ СДЕЛАНО|НУЖЕН ЧЕЛОВЕК|НЕ СМОГ|УСТАРЕЛО|ПУСТО)" + r"^ИТОГ:\s*(ГОТОВО|УЖЕ СДЕЛАНО|НУЖЕН ЧЕЛОВЕК|НЕ СМОГ|ПУСТО)" 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*)$" @@ -88,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: @@ -412,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 @@ -436,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 "" @@ -449,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 @@ -482,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) @@ -543,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``. @@ -577,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 " @@ -826,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__": @@ -882,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/pipeline_insights.py b/promptpilot/pipeline_insights.py index 8227cc5..0c9497d 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 @@ -1429,8 +1445,9 @@ def _adaptive_cadence_status(queue: dict, matching: dict | None, } -def _reconcile_adaptive_cadence(queue: dict, matching: dict | None, - backlog: int | None) -> dict | 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): @@ -1444,6 +1461,7 @@ def _reconcile_adaptive_cadence(queue: dict, matching: dict | None, 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({ @@ -1571,10 +1589,24 @@ def _window_metrics(snapshots: list[dict], current: dict, series_ids: list[int], 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) + 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()) - measured_hours = min(coverage, hours) + # 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 @@ -1584,15 +1616,16 @@ def hourly(value: int) -> float | None: "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)), - "exited_per_hour": hourly(len(exited)), + "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) for queue_id, count in queue_throughput.items()}, + 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)), } @@ -1692,6 +1725,8 @@ def worker_lane_policy() -> dict | None: 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 @@ -3334,8 +3369,28 @@ def _analyze_without_budget(profile_id: str, series: list[dict], *, # 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. - for item, matching, backlog in cadence_reconciliations: - _reconcile_adaptive_cadence(item, matching, backlog) + 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/worker.py b/promptpilot/worker.py index 92670ca..e54e0c6 100644 --- a/promptpilot/worker.py +++ b/promptpilot/worker.py @@ -136,8 +136,7 @@ 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 = ( - "ГОТОВО", "УЖЕ СДЕЛАНО", "НУЖЕН ЧЕЛОВЕК", "НЕ СМОГ", "УСТАРЕЛО", - "ПУСТО", + "ГОТОВО", "УЖЕ СДЕЛАНО", "НУЖЕН ЧЕЛОВЕК", "НЕ СМОГ", "ПУСТО", ) VERDICT_RE = re.compile(r"^[ \t>*#-]*ИТОГ:\s*(" + "|".join(VERDICTS) + r")\b", re.M | re.I) @@ -486,20 +485,11 @@ 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 - # A stale exact target is neither a failed run nor useful work. Re-elect it - # immediately instead of making a busy queue wait for its ordinary cadence. - # Only the explicit closing verdict enables this path: generic failures and - # ambiguous gate errors keep their normal bounded schedule. - stale_pipeline_target = False - if (task.verdict or "").upper() == "УСТАРЕЛО": - try: - from . import pipeline_insights - stale_pipeline_target = pipeline_insights._matching_queue(task) is not None - except (AttributeError, OSError, TypeError, ValueError): - # A broken/temporarily unreadable optional profile must not turn - # an ordinary recurring task into an unbounded immediate loop. - stale_pipeline_target = False - next_dt = (datetime.now(timezone.utc) if stale_pipeline_target + # 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 @@ -584,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). @@ -645,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) @@ -682,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( @@ -782,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 @@ -852,9 +847,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: " @@ -887,6 +891,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 @@ -1128,7 +1133,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( @@ -1426,7 +1432,8 @@ def _claim_next_task(busy_keys=(), busy_lane_ids=()): try: policy = pipeline_insights.worker_lane_policy() - except (OSError, TypeError, ValueError, json.JSONDecodeError) as exc: + except (AttributeError, OSError, TypeError, ValueError, + json.JSONDecodeError) as exc: print(f" !! pipeline lane scheduler unavailable: {exc}", flush=True) policy = None if policy is None: diff --git a/tests/test_herdr_workflow_completion.py b/tests/test_herdr_workflow_completion.py index 901661b..1c1a0df 100644 --- a/tests/test_herdr_workflow_completion.py +++ b/tests/test_herdr_workflow_completion.py @@ -37,11 +37,20 @@ def test_closing_workflow_verdict_accepts_only_final_line(): ) == "ГОТОВО" -def test_closing_workflow_verdict_accepts_stale_reselection_outcome(): +def test_stale_reselection_requires_targeted_exact_gate_fallback_form(): assert _closing_workflow_verdict( "Gate proved that the exact HEAD changed.\n" - "ИТОГ: УСТАРЕЛО (PR HEAD changed before mutation)" + "ИТОГ: УСТАРЕЛО (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ИТОГ: ПУСТО (очередь пуста)" ) == "ПУСТО" @@ -61,6 +70,15 @@ def test_closing_workflow_verdict_accepts_stale_reselection_outcome(): ) == "ГОТОВО" +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 d992fb0..ef40834 100644 --- a/tests/test_schedule_series.py +++ b/tests/test_schedule_series.py @@ -70,9 +70,7 @@ def test_fresh_install_has_no_project_specific_pipeline_profiles(tmp_path, monke assert pipeline_insights.list_profiles() == [] -def test_stale_pipeline_verdict_reselects_immediately(isolated_db, monkeypatch): - monkeypatch.setattr( - pipeline_insights, "_profiles", lambda: {"example": PIPELINE_PROFILE}) +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() @@ -88,22 +86,35 @@ def test_stale_pipeline_verdict_reselects_immediately(isolated_db, monkeypatch): before + timedelta(seconds=2) -def test_stale_verdict_keeps_normal_cadence_outside_pipeline( - isolated_db, monkeypatch): - monkeypatch.setattr(pipeline_insights, "_profiles", lambda: {}) +def test_second_consecutive_stale_uses_normal_cadence_and_non_stale_resets( + isolated_db): task = isolated_db.create_task(TaskCreate( - prompt="Ordinary recurring task", recurrence="4h")) - claimed = isolated_db.get_next_runnable() + prompt="ExampleProject - REVIEW", recurrence="4h")) + first = isolated_db.get_next_runnable() isolated_db.mark_completed( - claimed.id, "ИТОГ: УСТАРЕЛО", verdict="УСТАРЕЛО") + 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(claimed) + 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): @@ -144,18 +155,139 @@ def test_adaptive_fix_cadence_returns_to_idle_at_threshold(isolated_db): 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 = [ @@ -175,6 +307,60 @@ def test_window_metrics_report_actual_hourly_rates(isolated_db, monkeypatch): 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 = [] @@ -985,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 = { @@ -1581,9 +1802,9 @@ def run_analysis(): real_reconcile = pipeline_insights._reconcile_adaptive_cadence - def track_reconcile(*args): - cadence_calls.append(args) - return real_reconcile(*args) + 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) @@ -1615,6 +1836,53 @@ def track_reconcile(*args): 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 = { diff --git a/tests/test_worker_admission.py b/tests/test_worker_admission.py index 108c9c2..ee801c7 100644 --- a/tests/test_worker_admission.py +++ b/tests/test_worker_admission.py @@ -99,6 +99,22 @@ def test_lane_claim_is_atomic_and_queue_order_beats_fifo( 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 = {} @@ -172,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 @@ -183,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 = [] From 206921b70794c8337ac5ef08cce3018fc3b24c15 Mon Sep 17 00:00:00 2001 From: ibrog Date: Thu, 17 Sep 2026 11:58:01 +0300 Subject: [PATCH 3/3] =?UTF-8?q?fix(pipeline):=20=D0=BE=D0=B3=D1=80=D0=B0?= =?UTF-8?q?=D0=BD=D0=B8=D1=87=D0=B8=D1=82=D1=8C=20stale=20=D1=86=D0=B5?= =?UTF-8?q?=D0=BB=D0=B5=D0=B2=D1=8B=D0=BC=20fallback?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Generated-with: Codex --- promptpilot/pipeline_insights.py | 9 ++++- promptpilot/worker.py | 7 +++- tests/test_schedule_series.py | 60 ++++++++++++++++++++++++++++++++ 3 files changed, 74 insertions(+), 2 deletions(-) diff --git a/promptpilot/pipeline_insights.py b/promptpilot/pipeline_insights.py index 0c9497d..fee1d7d 100644 --- a/promptpilot/pipeline_insights.py +++ b/promptpilot/pipeline_insights.py @@ -2333,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, } diff --git a/promptpilot/worker.py b/promptpilot/worker.py index e54e0c6..41a8b10 100644 --- a/promptpilot/worker.py +++ b/promptpilot/worker.py @@ -830,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, diff --git a/tests/test_schedule_series.py b/tests/test_schedule_series.py index ef40834..dac4e07 100644 --- a/tests/test_schedule_series.py +++ b/tests/test_schedule_series.py @@ -2851,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", @@ -2877,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(