From 06e1fa6bfc92fb5dc85128b11c7ef0ccd9262f17 Mon Sep 17 00:00:00 2001 From: ibrog Date: Thu, 17 Sep 2026 22:59:27 +0300 Subject: [PATCH] fix(pipeline): bound GitHub budget starvation Give the oldest durable budget waiter an opt-in FIFO admission baton and a matching worker-slot preference without changing its quota floor. Preserve exact-attempt fencing across retries, restarts, migrations, configuration changes, and non-reserving execution paths. Generated-with: Codex --- README.md | 25 ++ promptpilot/db.py | 315 +++++++++++++++++--- promptpilot/pipeline_insights.py | 180 +++++++++++- promptpilot/worker.py | 17 +- tests/test_github_budget.py | 488 ++++++++++++++++++++++++++++++- 5 files changed, 968 insertions(+), 57 deletions(-) diff --git a/README.md b/README.md index d6fda8d..de45eea 100644 --- a/README.md +++ b/README.md @@ -795,6 +795,7 @@ Worker публикует heartbeat в общей SQLite БД. Активный "reset_grace_seconds": 60, "lease_seconds": 900, "busy_retry_seconds": 30, + "starvation_timeout_seconds": 900, "unavailable_retry_seconds": 300 }, "priority_control": { @@ -909,6 +910,16 @@ claim одной задачи происходят в одной SQLite-тран получат одно вхождение. Без `scheduler` остаётся прежний порядок `priority → created_at → id`. +Полоса становится выделенным физическим слотом только когда +`PP_CONCURRENCY` не меньше числа настроенных `lanes`. При меньшей +конкурентности это остаются детерминированные предпочтения свободных слотов, а +не резерв ёмкости; отдельную гарантию от ожидания GitHub API даёт описанный ниже +`github_budget.starvation_timeout_seconds`. +При включённом starvation timeout конфигурации `scheduler` нужна хотя бы одна +recovery-полоса с `borrow: true`; иначе lane policy считается невалидной и worker +откатывается к глобальному claim, чтобы переименованная очередь не удерживала +admission-эстафету навсегда. + #### Параллельные реплики REVIEW Очередь REVIEW может явно владеть несколькими независимыми сериями. Для этого @@ -1081,6 +1092,20 @@ skill остаётся дорогим: стартовые `3750 + 250` Core со аварии процесса. Приоритетная эстафета сохраняется до фактического резервирования бюджета, поэтому уже запущенный менее приоритетный этап не обгоняет разбуженный. Жёсткий reset живого лимита этим сигналом не сокращается. +Необязательный `starvation_timeout_seconds` (по умолчанию `0`, выключено; +рекомендуемое начальное значение `900`) ограничивает голодание при непрерывном +потоке более приоритетных запусков. После дедлайна самая старая задача с +durable reservation/priority/scan handoff (но не с ожиданием живого quota reset) +получает временную admission-эстафету: +новые резервы ждут, пока текущие освободятся и она зарезервирует бюджет. Более +новая задача не обходит более старую; при одинаковом времени ожидания порядок +задают priority и id. Эстафета меняет только порядок: исходный priority задачи, +`minimum_remaining` и `priority_one_headroom` не меняются, поэтому P3/P4 не +заимствуют запас P1. Если без чужих reservations живого лимита всё равно мало, +задача переходит в обычное ожидание reset и сразу освобождает эстафету. Для +профилей с общим GitHub token действует минимальный положительный timeout; +отсутствие настройки во всех профилях полностью сохраняет прежнюю строгую +приоритетную семантику. Недоступный лимит, повреждённый ledger или некорректная конфигурация работают fail-closed на `unavailable_retry_seconds`. Для профилей с одним GitHub token и `PP_DATA_DIR` PromptPilot покомпонентно применяет самый большой настроенный hard diff --git a/promptpilot/db.py b/promptpilot/db.py index 69f7355..9aabc2b 100644 --- a/promptpilot/db.py +++ b/promptpilot/db.py @@ -72,6 +72,7 @@ worktree_branch TEXT, note TEXT, budget_wait_scope TEXT, + budget_wait_started_at TEXT, pipeline_priority_restore INTEGER, verdict TEXT ,series_id INTEGER REFERENCES task_series(id) @@ -395,6 +396,8 @@ "ALTER TABLE tasks ADD COLUMN budget_wait_scope TEXT", "CREATE INDEX IF NOT EXISTS idx_tasks_budget_wait ON tasks(budget_wait_scope, status, priority)", "ALTER TABLE tasks ADD COLUMN pipeline_priority_restore INTEGER", + "ALTER TABLE tasks ADD COLUMN budget_wait_started_at TEXT", + "CREATE INDEX IF NOT EXISTS idx_tasks_budget_wait_fairness ON tasks(budget_wait_scope, status, budget_wait_started_at, priority, id)", ] WORKFLOW_SCHEMA_VERSION = "workflow_orchestrator_w0_v1" @@ -403,6 +406,8 @@ PIPELINE_TARGET_ATTEMPT_SCHEMA_VERSION = "pipeline_target_reservations_attempt_v2" PIPELINE_TARGET_OWNERSHIP_SCHEMA_VERSION = "pipeline_target_reservations_ownership_v3" PIPELINE_TARGET_HERDR_SCHEMA_VERSION = "pipeline_target_reservations_herdr_v4" +PIPELINE_BUDGET_WAITER_FAIRNESS_SCHEMA_VERSION = \ + "pipeline_github_budget_waiter_fairness_v1" def _now() -> str: @@ -434,6 +439,7 @@ def _row_to_task(row: sqlite3.Row) -> TaskInDB: # model/API. It identifies interruptible budget waits and a temporary # route-derived priority that must not leak into recurring successors. d.pop("budget_wait_scope", None) + d.pop("budget_wait_started_at", None) d.pop("pipeline_priority_restore", None) for field in ("scheduled_at", "next_run_at", "created_at", "started_at", "completed_at"): d[field] = _parse_dt(d[field]) @@ -516,6 +522,28 @@ def _init_db_once(): VALUES (?, ?)""", (PIPELINE_TARGET_HERDR_SCHEMA_VERSION, _now()), ) + fairness_migrated = conn.execute( + "SELECT 1 FROM schema_migrations WHERE version = ?", + (PIPELINE_BUDGET_WAITER_FAIRNESS_SCHEMA_VERSION,), + ).fetchone() + if fairness_migrated is None: + migration_now = _now() + # Existing durable reservation handoffs predate the age column. + # Give active legacy waiters one common baseline instead of + # leaving them permanently ineligible for the fairness baton. + conn.execute( + """UPDATE tasks SET budget_wait_started_at = ? + WHERE budget_wait_scope IS NOT NULL + AND budget_wait_started_at IS NULL + AND status IN ('pending', 'running')""", + (migration_now,), + ) + conn.execute( + """INSERT INTO schema_migrations (version, applied_at) + VALUES (?, ?)""", + (PIPELINE_BUDGET_WAITER_FAIRNESS_SCHEMA_VERSION, + migration_now), + ) _backfill_task_series(conn) @@ -1118,8 +1146,11 @@ 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, - order_key_fn=None) -> Optional[TaskInDB]: +def get_next_runnable( + busy_keys=(), key_fn=None, order_key_fn=None, *, + budget_wait_scope: str | None = None, + budget_starvation_timeout_seconds: int = 0, + budget_fairness_now: float | None = 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 @@ -1136,12 +1167,39 @@ def get_next_runnable(busy_keys=(), key_fn=None, now = _now() busy = set(busy_keys or ()) with _connect(immediate=True) as conn: + # The admission baton must reach a worker slot too. Otherwise an + # endless stream of higher-priority runnable tasks can repeatedly be + # claimed and denied while the selected waiter never reaches preflight. + fairness_task_id = None + if budget_starvation_timeout_seconds: + if (not isinstance(budget_wait_scope, str) + or not budget_wait_scope.strip()): + raise ValueError( + "budget wait scope is required for scheduler fairness") + fairness = _pipeline_github_budget_waiter_summary( + conn, budget_wait_scope, + starvation_timeout_seconds=( + budget_starvation_timeout_seconds), + now=budget_fairness_now).get("fairness_waiter") + if isinstance(fairness, dict): + fairness_task_id = fairness.get("task_id") + # A reservation wait normally carries a short retry deadline. + # Once its baton matures, make the selected pending attempt + # runnable in this same claim transaction; otherwise it could + # block new admissions while sleeping until that deadline. + conn.execute( + """UPDATE tasks SET scheduled_at = ?, next_run_at = NULL + WHERE id = ? AND status = 'pending' + AND budget_wait_scope = ?""", + (now, fairness_task_id, budget_wait_scope), + ) # Sequential worker takes just the top task. When some keys are busy we # walk the whole runnable queue in priority order until a non-colliding # 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 or order_key_fn is not None else " LIMIT 1" + limit_clause = "" if (busy or order_key_fn is not None + or fairness_task_id is not None) else " LIMIT 1" rows = conn.execute( f"""SELECT * FROM tasks WHERE status IN ('pending', 'rate_limited') @@ -1158,11 +1216,12 @@ def get_next_runnable(busy_keys=(), key_fn=None, 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: + policy_rank = order_key_fn(task) if order_key_fn is not None else () + if order_key_fn is not None and policy_rank is None: continue - candidates.append((rank, position, task)) - if order_key_fn is not None: + fairness_rank = 0 if task.id == fairness_task_id else 1 + candidates.append(((fairness_rank, policy_rank), position, task)) + if order_key_fn is not None or fairness_task_id is not None: candidates.sort(key=lambda item: (item[0], item[1])) for _rank, _position, task in candidates: started_at = _now() @@ -1278,7 +1337,8 @@ def mark_completed(task_id: int, result: str, exit_code: int = 0, "next_run_at = NULL, exit_code = ?, completed_at = ?, model_used = ?, " "session_id = COALESCE(?, session_id), " f"verdict = COALESCE(?, verdict), note = NULL, " - f"budget_wait_scope = NULL WHERE {where}", + f"budget_wait_scope = NULL, budget_wait_started_at = NULL " + f"WHERE {where}", values, ) if cur.rowcount: @@ -1298,7 +1358,8 @@ def mark_failed(task_id: int, error: str, exit_code: int = 1, values.append(attempt) cur = conn.execute( "UPDATE tasks SET status = 'failed', error = ?, exit_code = ?, " - f"completed_at = ?, note = NULL, budget_wait_scope = NULL " + f"completed_at = ?, note = NULL, budget_wait_scope = NULL, " + f"budget_wait_started_at = NULL " f"WHERE {where}", values, ) @@ -1324,7 +1385,8 @@ def fail_running_attempt(task_id: int, started_at, error: str, cur = conn.execute( """UPDATE tasks SET status = 'failed', error = ?, exit_code = ?, - completed_at = ?, note = NULL, budget_wait_scope = NULL + completed_at = ?, note = NULL, budget_wait_scope = NULL, + budget_wait_started_at = NULL WHERE id = ? AND status = 'running' AND started_at = ?""", (error, exit_code, _now(), task_id, started_at), ) @@ -1348,7 +1410,8 @@ def mark_rate_limited(task_id: int, next_run_at: datetime, error: str = None, SET status = 'rate_limited', next_run_at = ?, retry_count = retry_count + 1, - error = COALESCE(?, error), budget_wait_scope = NULL + error = COALESCE(?, error), budget_wait_scope = NULL, + budget_wait_started_at = NULL WHERE """ + where, values, ) @@ -1404,6 +1467,7 @@ def defer_task(task_id: int, next_run_at: datetime, reason: str = None, attempt = _attempt_iso(expected_started_at) where = "id = ? AND status = 'running'" values = [deadline, deadline if hard_not_before else None, reason, + budget_wait_scope, budget_wait_scope, _now(), budget_wait_scope, task_id] if attempt is not None: where += " AND started_at = ?" @@ -1411,7 +1475,15 @@ def defer_task(task_id: int, next_run_at: datetime, reason: str = None, cur = conn.execute( """UPDATE tasks SET status = 'pending', scheduled_at = ?, next_run_at = ?, started_at = NULL, - error = COALESCE(?, error), budget_wait_scope = ? + error = COALESCE(?, error), + budget_wait_started_at = CASE + WHEN ? IS NULL THEN NULL + WHEN budget_wait_scope = ? + AND budget_wait_started_at IS NOT NULL + THEN budget_wait_started_at + ELSE ? + END, + budget_wait_scope = ? WHERE """ + where, values, ) @@ -1457,7 +1529,8 @@ def mark_cancelled(task_id: int, note: str = None, cur = conn.execute( "UPDATE tasks SET status = 'cancelled', completed_at = ?, " f"error = COALESCE(?, error), note = NULL, " - f"budget_wait_scope = NULL WHERE {where}", + f"budget_wait_scope = NULL, budget_wait_started_at = NULL " + f"WHERE {where}", values, ) if cur.rowcount: @@ -1470,7 +1543,8 @@ def cancel_task(task_id: int) -> bool: with _connect() as conn: cur = conn.execute( "UPDATE tasks SET status = 'cancelled', completed_at = ?, " - "budget_wait_scope = NULL WHERE id = ? " + "budget_wait_scope = NULL, budget_wait_started_at = NULL " + "WHERE id = ? " "AND status IN ('pending', 'rate_limited')", (_now(), task_id), ) @@ -1482,9 +1556,13 @@ def cancel_task(task_id: int) -> bool: def update_priority(task_id: int, priority: int) -> bool: with _connect() as conn: cur = conn.execute( - """UPDATE tasks SET priority = ?, pipeline_priority_restore = NULL + """UPDATE tasks SET priority = ?, pipeline_priority_restore = NULL, + budget_wait_started_at = CASE + WHEN budget_wait_scope IS NOT NULL THEN ? + ELSE NULL + END WHERE id = ? AND status IN ('pending', 'rate_limited')""", - (priority, task_id), + (priority, _now(), task_id), ) return cur.rowcount > 0 @@ -1730,8 +1808,14 @@ def update_series(series_id: int, fields: dict) -> bool: task_values = list(task_fields.values()) if "priority" in fields: # An explicit series edit supersedes a temporary route-derived - # promotion on the waiting occurrence. - task_sets.append("pipeline_priority_restore = NULL") + # promotion and restarts any admission-wait ordering for the + # edited occurrence. + task_sets.extend([ + "pipeline_priority_restore = NULL", + "budget_wait_started_at = CASE " + "WHEN budget_wait_scope IS NOT NULL THEN ? ELSE NULL END", + ]) + task_values.append(_now()) if provider_changed: # Resume ids and retry budgets belong to the old provider. Carrying # them across a switch can ask a new CLI to resume an incompatible @@ -2098,6 +2182,14 @@ def series_action(series_id: int, action: str) -> bool: if action == "pause": conn.execute("UPDATE task_series SET paused = 1, updated_at = ? WHERE id = ?", (now, series_id)) + conn.execute( + """UPDATE tasks + SET budget_wait_scope = NULL, + budget_wait_started_at = NULL + WHERE series_id = ? + AND status IN ('pending', 'rate_limited')""", + (series_id,), + ) conn.execute( "DELETE FROM settings WHERE key = ?", (_pipeline_series_wake_intent_key(series_id),), @@ -2107,7 +2199,10 @@ def series_action(series_id: int, action: str) -> bool: (now, series_id)) elif action == "run_now": cur = conn.execute( - """UPDATE tasks SET scheduled_at = ?, next_run_at = NULL + """UPDATE tasks + SET scheduled_at = ?, next_run_at = NULL, + budget_wait_scope = NULL, + budget_wait_started_at = NULL WHERE series_id = ? AND status IN ('pending', 'rate_limited')""", (now, series_id), ) @@ -2141,7 +2236,8 @@ def series_action(series_id: int, action: str) -> bool: conn.execute("UPDATE task_series SET ended_at = ?, updated_at = ? WHERE id = ?", (now, now, series_id)) conn.execute("UPDATE tasks SET status = 'cancelled', completed_at = ?, " - "budget_wait_scope = NULL " + "budget_wait_scope = NULL, " + "budget_wait_started_at = NULL " "WHERE series_id = ? AND status IN ('pending', 'rate_limited')", (now, series_id)) conn.execute( @@ -2420,21 +2516,31 @@ def update_task_fields(task_id: int, fields: dict) -> bool: fields = {k: v for k, v in fields.items() if k in EDITABLE_FIELDS} if not fields: return False - if "scheduled_at" in fields: + explicit_reschedule = "scheduled_at" in fields + if explicit_reschedule: fields["scheduled_at"] = _to_utc_iso(fields["scheduled_at"]) # An explicit operator reschedule supersedes an automatic soft wait. # Without this, a task intentionally moved far into the future would # keep lower-priority work from reserving GitHub budget indefinitely. fields["budget_wait_scope"] = None + fields["budget_wait_started_at"] = None if "priority" in fields: # Do not let a later route restoration undo the operator's explicit - # priority choice for a temporarily promoted waiting occurrence. + # priority choice for a temporarily promoted waiting occurrence. The + # edit also starts a fresh fairness history for the new priority. fields["pipeline_priority_restore"] = None - sets = ", ".join(f"{k} = ?" for k in fields) + sets = [f"{k} = ?" for k in fields] + values = list(fields.values()) + if "priority" in fields and not explicit_reschedule: + sets.append( + "budget_wait_started_at = CASE " + "WHEN budget_wait_scope IS NOT NULL THEN ? ELSE NULL END") + values.append(_now()) with _connect() as conn: cur = conn.execute( - f"UPDATE tasks SET {sets} WHERE id = ? AND status IN ('pending', 'rate_limited')", - (*fields.values(), task_id), + f"UPDATE tasks SET {', '.join(sets)} " + "WHERE id = ? AND status IN ('pending', 'rate_limited')", + (*values, task_id), ) return cur.rowcount > 0 @@ -3011,7 +3117,19 @@ def _write_pipeline_budget_reservations(conn, key: str, values: list[dict]) -> N ) -def _pipeline_github_budget_waiter_summary(conn, scope: str) -> dict: +def _pipeline_budget_starvation_timeout(value: int) -> int: + if (isinstance(value, bool) or not isinstance(value, int) + or not 0 <= value <= 86400): + raise ValueError( + "pipeline budget starvation timeout must be 0..86400 seconds") + return value + + +def _pipeline_github_budget_waiter_summary( + conn, scope: str, *, starvation_timeout_seconds: int = 0, + now: Optional[float] = None) -> dict: + starvation_timeout_seconds = _pipeline_budget_starvation_timeout( + starvation_timeout_seconds) row = conn.execute( """SELECT COUNT(*) AS count, MIN(priority) AS min_priority FROM tasks @@ -3023,9 +3141,50 @@ def _pipeline_github_budget_waiter_summary(conn, scope: str) -> dict: AND s.paused = 0 AND s.ended_at IS NULL))))""", (scope,), ).fetchone() + fairness_waiter = None + if starvation_timeout_seconds: + oldest = conn.execute( + """SELECT id, priority, budget_wait_started_at + FROM tasks + WHERE budget_wait_scope = ? + AND budget_wait_started_at IS NOT NULL + AND (status = 'running' OR ( + status = 'pending' + AND (series_id IS NULL OR EXISTS ( + SELECT 1 FROM task_series s WHERE s.id = tasks.series_id + AND s.paused = 0 AND s.ended_at IS NULL)))) + ORDER BY budget_wait_started_at ASC, priority ASC, id ASC + LIMIT 1""", + (scope,), + ).fetchone() + if oldest is not None: + try: + wait_started = datetime.fromisoformat( + oldest["budget_wait_started_at"]) + if wait_started.tzinfo is None: + raise ValueError("missing timezone") + wait_started = wait_started.astimezone(timezone.utc) + except (TypeError, ValueError) as exc: + raise ValueError( + "pipeline budget waiter timestamp is invalid") from exc + current = (time.time() if now is None else float(now)) + if not math.isfinite(current): + raise ValueError("pipeline budget fairness time must be finite") + eligible_at = wait_started.timestamp() + starvation_timeout_seconds + if current >= eligible_at: + fairness_waiter = { + "task_id": int(oldest["id"]), + "priority": int(oldest["priority"]), + "wait_started_at": wait_started.isoformat(), + "eligible_at": datetime.fromtimestamp( + eligible_at, timezone.utc).isoformat(), + "wait_seconds": max( + 0, int(current - wait_started.timestamp())), + } return { "count": int(row["count"] or 0) if row is not None else 0, "min_priority": row["min_priority"] if row is not None else None, + "fairness_waiter": fairness_waiter, } @@ -3046,7 +3205,8 @@ def _wake_pipeline_github_budget_waiters(conn, scope: str) -> dict: def arm_pipeline_github_budget_waiter( scope: str, *, task_id: int, task_started_at, - expected_revision: int | None = None) -> dict: + expected_revision: int | None = None, + now: Optional[float] = None) -> dict: """Persist a priority handoff before releasing the admission scan lease. New reservations are serialized by that lease, while reservation release @@ -3067,12 +3227,23 @@ def arm_pipeline_github_budget_waiter( attempt = _attempt_iso(task_started_at) if attempt is None: raise ValueError("running task attempt is required for budget waiter") + current = time.time() if now is None else float(now) + if not math.isfinite(current): + raise ValueError("pipeline budget waiter time must be finite") + wait_started_at = datetime.fromtimestamp(current, timezone.utc).isoformat() with _connect(immediate=True) as conn: revision = _pipeline_github_budget_revision(conn, scope) cur = conn.execute( - """UPDATE tasks SET budget_wait_scope = ? + """UPDATE tasks + SET budget_wait_started_at = CASE + WHEN budget_wait_scope = ? + AND budget_wait_started_at IS NOT NULL + THEN budget_wait_started_at + ELSE ? + END, + budget_wait_scope = ? WHERE id = ? AND status = 'running' AND started_at = ?""", - (scope, task_id, attempt), + (scope, wait_started_at, scope, task_id, attempt), ) return { "armed": cur.rowcount > 0, @@ -3083,7 +3254,27 @@ def arm_pipeline_github_budget_waiter( } -def pipeline_github_budget_reservations(scope: str) -> dict: +def clear_pipeline_github_budget_waiter( + *, task_id: int, task_started_at) -> bool: + """Clear one successful attempt's handoff without touching a successor.""" + if isinstance(task_id, bool) or not isinstance(task_id, int) or task_id <= 0: + raise ValueError("running task id is required for budget waiter clear") + attempt = _attempt_iso(task_started_at) + if attempt is None: + raise ValueError("running task attempt is required for budget waiter clear") + with _connect() as conn: + cur = conn.execute( + """UPDATE tasks + SET budget_wait_scope = NULL, budget_wait_started_at = NULL + WHERE id = ? AND status = 'running' AND started_at = ?""", + (task_id, attempt), + ) + return cur.rowcount > 0 + + +def pipeline_github_budget_reservations( + scope: str, *, starvation_timeout_seconds: int = 0, + now: Optional[float] = None) -> dict: """Return live, task-fenced reservations and prune completed attempts.""" key = _pipeline_github_budget_reservation_key(scope) with _connect(immediate=True) as conn: @@ -3111,10 +3302,17 @@ def pipeline_github_budget_reservations(scope: str) -> dict: for item in active: for resource in totals: totals[resource] += item["cost"][resource] - waiting = _pipeline_github_budget_waiter_summary(conn, scope) - return {"count": len(active), "totals": totals, "items": active, - "revision": revision, "woken_waiters": woken, - "waiting_waiters": waiting} + waiting = _pipeline_github_budget_waiter_summary( + conn, scope, + starvation_timeout_seconds=starvation_timeout_seconds, + now=now) + fairness_waiter = waiting.pop("fairness_waiter", None) + result = {"count": len(active), "totals": totals, "items": active, + "revision": revision, "woken_waiters": woken, + "waiting_waiters": waiting} + if fairness_waiter is not None: + result["fairness_waiter"] = fairness_waiter + return result def record_pipeline_github_rate_snapshot( @@ -3205,6 +3403,7 @@ def reserve_pipeline_github_budget( scope: str, *, token: str, task_id: int, task_started_at: str, profile_id: str, queue_id: str, route: str, cost: dict, limits: dict, minimum_remaining: dict, + starvation_timeout_seconds: int = 0, now: Optional[float] = None, scan_lease_guard: Optional[dict] = None, expected_revision: Optional[int] = None) -> dict: @@ -3228,6 +3427,8 @@ def reserve_pipeline_github_budget( "expected pipeline budget revision must be a non-negative integer") requested = _pipeline_budget_vector(cost, "cost") floor = _pipeline_budget_vector(minimum_remaining, "minimum_remaining") + starvation_timeout_seconds = _pipeline_budget_starvation_timeout( + starvation_timeout_seconds) reported = {} resets = {} for resource in _PIPELINE_GITHUB_BUDGET_RESOURCES: @@ -3294,7 +3495,9 @@ def reserve_pipeline_github_budget( or item["cost"] != requested or item["route"] != route): raise ValueError("pipeline budget token was reused inconsistently") conn.execute( - """UPDATE tasks SET budget_wait_scope = NULL + """UPDATE tasks + SET budget_wait_scope = NULL, + budget_wait_started_at = NULL WHERE id = ? AND status = 'running' AND started_at = ? AND budget_wait_scope = ?""", (task_id, task_started_at, scope), @@ -3333,9 +3536,27 @@ def reserve_pipeline_github_budget( "revision": revision, } - waiting = _pipeline_github_budget_waiter_summary(conn, scope) + waiting = _pipeline_github_budget_waiter_summary( + conn, scope, + starvation_timeout_seconds=starvation_timeout_seconds, + now=created_at) + fairness_waiter = waiting.get("fairness_waiter") + if (isinstance(fairness_waiter, dict) + and fairness_waiter.get("task_id") != task_id): + return { + "allowed": False, "state": "fairness_waiter", + "reason": ( + "the oldest GitHub budget waiter owns the admission baton"), + "fairness_waiter": fairness_waiter, + "blocked_resources": [], + "reported_remaining": reported, "reserved_other": reserved, + "requested_cost": requested, "minimum_remaining": floor, + "effective_after": after, "active_reservations": len(active), + "revision": revision, + } waiting_priority = waiting.get("min_priority") - if (type(waiting_priority) is int + if (fairness_waiter is None + and type(waiting_priority) is int and waiting_priority < int(task["priority"])): return { "allowed": False, "state": "priority_waiter", @@ -3387,6 +3608,19 @@ def reserve_pipeline_github_budget( state = ("low" if any(item["blocked_by"] == "live" for item in blocked) else "budget_in_flight") + if state == "low": + # A waiter that cannot fit under its unchanged quota floor + # must relinquish the baton atomically with this decision. A + # worker crash before defer must not freeze new reservations + # until GitHub resets. + conn.execute( + """UPDATE tasks + SET budget_wait_scope = NULL, + budget_wait_started_at = NULL + WHERE id = ? AND status = 'running' + AND started_at = ? AND budget_wait_scope = ?""", + (task_id, task_started_at, scope), + ) return {"allowed": False, "state": state, "blocked_resources": blocked, "reported_remaining": reported, "reserved_other": reserved, @@ -3403,7 +3637,9 @@ def reserve_pipeline_github_budget( _write_pipeline_budget_reservations(conn, key, active) revision = _advance_pipeline_github_budget_revision(conn, scope) conn.execute( - """UPDATE tasks SET budget_wait_scope = NULL + """UPDATE tasks + SET budget_wait_scope = NULL, + budget_wait_started_at = NULL WHERE id = ? AND status = 'running' AND started_at = ? AND budget_wait_scope = ?""", (task_id, task_started_at, scope), @@ -4005,7 +4241,8 @@ def recover_running_attempt(task_id: int, started_at) -> bool: """UPDATE tasks SET status = 'cancelled', completed_at = ?, error = COALESCE(error, ?), note = NULL, - budget_wait_scope = NULL + budget_wait_scope = NULL, + budget_wait_started_at = NULL WHERE id = ? AND status = 'running' AND started_at = ?""", (_now(), "Отменена пользователем до перезапуска worker", task_id, attempt), diff --git a/promptpilot/pipeline_insights.py b/promptpilot/pipeline_insights.py index 4a1fb5d..02ee72b 100644 --- a/promptpilot/pipeline_insights.py +++ b/promptpilot/pipeline_insights.py @@ -362,6 +362,12 @@ def _github_budget_policy(profile: dict) -> dict | None: raw, "busy_retry_seconds", 30, 5, 300), "unavailable_retry_seconds": _bounded_int( raw, "unavailable_retry_seconds", 300, 30, 3600), + # Opt in explicitly: once this deadline expires the oldest durable + # reservation waiter temporarily outranks new admissions, including + # priority 1. Its own priority floor still applies, so fairness drains + # competing reservations without borrowing urgent-task headroom. + "starvation_timeout_seconds": _bounded_int( + raw, "starvation_timeout_seconds", 0, 0, 86400), # GitHub primary budgets belong to the authenticated account, not to a # repository/profile. Keep one scope per shared PromptPilot database so # profiles cannot accidentally opt out of each other's reservation. @@ -373,6 +379,9 @@ def _with_shared_budget_floor(policy: dict) -> dict: """Use the strongest hard reserve of every profile sharing this account.""" floor = dict(policy["minimum_remaining"]) priority_one_headroom = dict(policy["priority_one_headroom"]) + starvation_timeouts = [] + if policy.get("starvation_timeout_seconds", 0) > 0: + starvation_timeouts.append(policy["starvation_timeout_seconds"]) for configured_profile in _profiles().values(): candidate = _github_budget_policy(configured_profile) if candidate is None or candidate["lease_scope"] != policy["lease_scope"]: @@ -382,9 +391,16 @@ def _with_shared_budget_floor(policy: dict) -> dict: for resource, value in candidate["priority_one_headroom"].items(): priority_one_headroom[resource] = max( priority_one_headroom[resource], value) + if candidate.get("starvation_timeout_seconds", 0) > 0: + starvation_timeouts.append( + candidate["starvation_timeout_seconds"]) selected = dict(policy) selected["minimum_remaining"] = floor selected["priority_one_headroom"] = priority_one_headroom + # The reservation ledger is account-wide. Every caller must therefore use + # one deterministic deadline; the tightest explicit bound wins. + selected["starvation_timeout_seconds"] = ( + min(starvation_timeouts) if starvation_timeouts else 0) if any(priority_one_headroom.values()) and selected.get("costs") is None: raise ValueError( "общий github_budget.priority_one_headroom требует " @@ -736,12 +752,49 @@ def _evaluate_github_budget(policy: dict, limits: dict | None, *, } +def _budget_reservation_ledger(policy: dict, *, now: float | None = None) -> dict: + """Read the shared ledger, adding fairness options only when enabled.""" + timeout = int(policy.get("starvation_timeout_seconds") or 0) + if not timeout: + return db.pipeline_github_budget_reservations(policy["lease_scope"]) + return db.pipeline_github_budget_reservations( + policy["lease_scope"], starvation_timeout_seconds=timeout, now=now) + + def _priority_waiter_decision( policy: dict, limits: dict | None, reservations: dict, *, - priority: int | None, status_revision: int | None) -> dict | None: - """Yield to an already waiting or just-woken higher-priority task.""" + priority: int | None, task_id: int | None, + status_revision: int | None) -> dict | None: + """Yield to priority normally, then to an aged oldest-waiter baton.""" woken = reservations.get("woken_waiters") or {} waiting = reservations.get("waiting_waiters") or {} + fairness_waiter = ( + reservations.get("fairness_waiter") + or waiting.get("fairness_waiter")) + if isinstance(fairness_waiter, dict): + fair_task_id = fairness_waiter.get("task_id") + if type(fair_task_id) is int and fair_task_id == task_id: + # Aging affects ordering only. The selected task proceeds with its + # unchanged admission priority and hard/headroom floor. + return None + now = time.time() + decision = _budget_denied( + policy, state="fairness_waiter", + reason=( + "GitHub API-бюджет передан самой старой ожидающей задаче " + f"#{fair_task_id} после {fairness_waiter.get('wait_seconds', 0)} с " + "ожидания; новые резервы временно приостановлены"), + now=now, limits=limits, + defer_at=now + policy["busy_retry_seconds"], + status_revision=status_revision, + reserved_other=reservations.get("totals"), + active_reservations=int(reservations.get("count") or 0), + ) + decision["fairness_waiter"] = copy.deepcopy(fairness_waiter) + revision = reservations.get("revision") + if type(revision) is int and revision >= 0: + decision["_budget_reservation_revision"] = revision + return decision priorities = [ value for value in ( woken.get("min_priority"), waiting.get("min_priority")) @@ -775,7 +828,7 @@ def _arm_budget_waiter(task, decision: dict) -> dict: state = decision.get("state") if state not in { "budget_in_flight", "ledger_changed", "priority_waiter", - "scan_in_progress"}: + "fairness_waiter", "scan_in_progress"}: return decision scope = decision.get("lease_scope") revision = decision.get("_budget_reservation_revision") @@ -831,11 +884,14 @@ def _github_scan_admission(profile: dict, purpose: str, "unavailable_retry_seconds": 300, "lease_scope": _GITHUB_SCAN_LEASE_SCOPE, } - yield _budget_denied( + decision = _budget_denied( fallback, state="invalid_config", reason=f"Некорректный github_budget: {exc}", now=now) + _clear_budget_waiter_after_denial(task) + yield decision return if policy is None: + _clear_budget_waiter_after_denial(task) yield {"enabled": False, "allowed": True, "state": "legacy"} return @@ -844,9 +900,11 @@ def _github_scan_admission(profile: dict, purpose: str, acquired = db.acquire_pipeline_scan_lease( policy["lease_scope"], token, policy["lease_seconds"], now=now) except Exception as exc: - yield _budget_denied( + decision = _budget_denied( policy, state="lease_unavailable", reason=f"SQLite lease GitHub-сканирования недоступна: {exc}", now=now) + _clear_budget_waiter_after_denial(task) + yield decision return if not acquired.get("acquired"): status_revision = acquired.get("status_revision") @@ -878,6 +936,10 @@ def _github_scan_admission(profile: dict, purpose: str, status_revision=( status_revision if isinstance(status_revision, int) else None)) + if decision.get("state") not in { + "budget_in_flight", "ledger_changed", "priority_waiter", + "fairness_waiter", "scan_in_progress"}: + _clear_budget_waiter_after_denial(task) yield decision return @@ -899,6 +961,8 @@ def _github_scan_admission(profile: dict, purpose: str, except _GitHubRateLimitUnavailable as exc: decision = _rate_limit_failure_decision( profile, exc, status_revision=lease.status_revision) + if task is not None: + _clear_budget_waiter_after_denial(task) decision["_lease"] = lease yield decision return @@ -909,8 +973,7 @@ def _github_scan_admission(profile: dict, purpose: str, status_revision=lease.status_revision) else: try: - reservations = db.pipeline_github_budget_reservations( - policy["lease_scope"]) + reservations = _budget_reservation_ledger(policy) except Exception as exc: decision = _budget_denied( policy, state="ledger_unavailable", @@ -924,6 +987,7 @@ def _github_scan_admission(profile: dict, purpose: str, decision = _priority_waiter_decision( policy, limits, reservations, priority=admission_priority, + task_id=getattr(task, "id", None), status_revision=lease.status_revision) if decision is None: decision = _evaluate_github_budget( @@ -942,6 +1006,11 @@ def _github_scan_admission(profile: dict, purpose: str, decision = _lease_failure_decision( profile, f"GitHub budget admission не выполнен: {exc}", status_revision=lease.status_revision, unavailable=True) + if (task is not None and not decision.get("allowed") + and decision.get("state") not in { + "budget_in_flight", "ledger_changed", "priority_waiter", + "fairness_waiter", "scan_in_progress"}): + _clear_budget_waiter_after_admission(task) decision["_lease"] = lease yield decision finally: @@ -1007,6 +1076,11 @@ def _reservation_denied(policy: dict, limits: dict | None, result: dict, *, state = "priority_waiter" reason = str(result.get("reason") or "GitHub budget yielded to a higher-priority waiter") + elif result.get("state") == "fairness_waiter": + defer_at = now + policy["busy_retry_seconds"] + state = "fairness_waiter" + reason = str(result.get("reason") or + "GitHub budget yielded to the oldest aged waiter") elif blocked: live_blocked = [item for item in blocked if item.get("blocked_by") == "live"] @@ -1040,9 +1114,36 @@ def _reservation_denied(policy: dict, limits: dict | None, result: dict, *, revision = result.get("revision") if type(revision) is int and revision >= 0: decision["_budget_reservation_revision"] = revision + if isinstance(result.get("fairness_waiter"), dict): + decision["fairness_waiter"] = copy.deepcopy( + result["fairness_waiter"]) return decision +def _clear_budget_waiter_after_admission(task) -> None: + """Settle a durable handoff when no in-flight reservation will do it.""" + task_id = getattr(task, "id", None) + started_at = getattr(task, "started_at", None) + # Synthetic/read-only callers cannot own a durable marker. Real worker + # attempts always carry both values and are fenced before provider launch. + if type(task_id) is not int or task_id <= 0 or started_at is None: + return + if not db.clear_pipeline_github_budget_waiter( + task_id=task_id, task_started_at=started_at): + raise ValueError( + "running task attempt changed before GitHub budget waiter clear") + + +def _clear_budget_waiter_after_denial(task) -> None: + """Best-effort cleanup for a denial that will not preserve a handoff.""" + try: + _clear_budget_waiter_after_admission(task) + except (sqlite3.Error, TypeError, ValueError): + # The denial remains fail-closed. If SQLite itself is unavailable the + # worker defer/terminal path will retry cleanup once storage recovers. + pass + + def _reserve_execution_admission( task, profile_id: str, profile: dict, queue: dict, admission: dict, budget_route: str, *, retain_budget: bool, @@ -1058,6 +1159,7 @@ def _reserve_execution_admission( try: policy = _github_budget_policy(profile) if policy is None: + _clear_budget_waiter_after_admission(task) return admission policy = _with_shared_budget_floor(policy) policy = _budget_policy_for_route(policy, budget_route) @@ -1066,6 +1168,7 @@ def _reserve_execution_admission( if not policy.get("cost_accounting"): # Profiles without route costs keep the original one-shot floor # admission; they neither re-read limits nor create reservations. + _clear_budget_waiter_after_admission(task) return admission if not isinstance(lease, _GitHubScanLease): raise _GitHubScanLeaseLost( @@ -1087,11 +1190,11 @@ def _reserve_execution_admission( policy, limits, now=time.time(), status_revision=lease.status_revision) elif not retain_budget: - reservations = db.pipeline_github_budget_reservations( - policy["lease_scope"]) + reservations = _budget_reservation_ledger(policy) decision = _priority_waiter_decision( policy, limits, reservations, priority=effective_priority, + task_id=getattr(task, "id", None), status_revision=lease.status_revision) if decision is None: decision = _evaluate_github_budget( @@ -1114,6 +1217,8 @@ def _reserve_execution_admission( queue_id=str(queue.get("id") or ""), route=budget_route, cost=policy["requested_cost"], limits=limits, minimum_remaining=policy["minimum_remaining"], + starvation_timeout_seconds=policy.get( + "starvation_timeout_seconds", 0), scan_lease_guard=lease.guard, expected_revision=expected_revision) if not reserved.get("allowed"): @@ -1146,6 +1251,8 @@ def _reserve_execution_admission( _GitHubBudgetReservation( policy["lease_scope"], token, task_id, started_at) decision = _arm_budget_waiter(task, decision) + if decision.get("allowed") and not retain_budget: + _clear_budget_waiter_after_admission(task) except _GitHubScanLeaseFailure as exc: decision = _lease_exception_decision( profile, exc, @@ -1168,6 +1275,10 @@ def _reserve_execution_admission( now=time.time(), status_revision=(lease.status_revision if isinstance(lease, _GitHubScanLease) else None)) + if (not decision.get("allowed") and decision.get("state") not in { + "budget_in_flight", "ledger_changed", "priority_waiter", + "fairness_waiter", "scan_in_progress"}): + _clear_budget_waiter_after_denial(task) decision["_lease"] = lease admission.clear() admission.update(decision) @@ -2254,9 +2365,45 @@ def worker_lane_policy() -> dict | None: raise ValueError( f"очередь {profile_id}/{queue_id} настроена на {replicas} " f"реплики, но доступна только в {eligible} scheduler lanes") + if lanes: + fairness_enabled = any( + (budget := _github_budget_policy(profile)) is not None + and budget.get("starvation_timeout_seconds", 0) > 0 + for profile in profiles.values() + ) + if fairness_enabled and not any(lane["borrow"] for lane in lanes): + raise ValueError( + "scheduler with GitHub budget fairness requires at least one " + "borrow:true recovery lane") return {"profiles": profiles, "lanes": lanes} if lanes else None +def worker_budget_fairness_policy() -> dict | None: + """Return the shared opt-in claim policy without touching GitHub. + + Admission fairness is ineffective if the selected task cannot win a local + worker slot. Every profile shares the authenticated account's ledger, so + the shortest positive configured timeout is also used by the scheduler. + """ + timeouts = [] + scope = None + for profile in _profiles().values(): + policy = _github_budget_policy(profile) + if policy is None or policy.get("starvation_timeout_seconds", 0) <= 0: + continue + timeouts.append(policy["starvation_timeout_seconds"]) + scope = policy["lease_scope"] if scope is None else scope + if scope != policy["lease_scope"]: + raise ValueError( + "GitHub budget fairness profiles must share one lease scope") + if not timeouts: + return None + return { + "scope": scope, + "starvation_timeout_seconds": min(timeouts), + } + + 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. @@ -2684,6 +2831,7 @@ def _budget_defer_route(admission: dict, profile_id: str, profile: dict, reservation_handoff = ( admission.get("state") in { "budget_in_flight", "ledger_changed", "priority_waiter", + "fairness_waiter", "scan_in_progress"} and type(reservation_revision) is int and reservation_revision >= 0 @@ -2694,7 +2842,8 @@ def _budget_defer_route(admission: dict, profile_id: str, profile: dict, reservation_handoff and admission.get("state") == "budget_in_flight") else "retry" if admission.get("state") in { - "ledger_changed", "priority_waiter", "scan_in_progress"} + "ledger_changed", "priority_waiter", "fairness_waiter", + "scan_in_progress"} else "hard_not_before") result = { "action": "defer", "mode": mode, @@ -2723,6 +2872,10 @@ def execution_route(task, fallback_prompt: str, working_dir: str | None = None, """ matched = _matching_queue(task) if matched is None: + # A profile/series marker may have been renamed while this exact + # attempt was waiting. It is now a legacy prompt path with no final + # budget reservation to settle the old shared-scope baton. + _clear_budget_waiter_after_admission(task) return {"action": "prompt", "mode": "skill", "prompt": fallback_prompt} profile_id, profile, queue = matched try: @@ -4001,8 +4154,7 @@ def _with_live_github_budget_state(result: dict, profile: dict) -> dict: reservations = None ledger_reason = None try: - reservations = db.pipeline_github_budget_reservations( - policy["lease_scope"]) + reservations = _budget_reservation_ledger(policy, now=now) except (OSError, sqlite3.Error, TypeError, ValueError) as exc: ledger_reason = str(exc) @@ -4057,6 +4209,10 @@ def _with_live_github_budget_state(result: dict, profile: dict) -> dict: route = str(item.get("route") or "unknown") route_counts[route] = route_counts.get(route, 0) + 1 budget["reservation_routes"] = route_counts + budget["waiting_waiters"] = copy.deepcopy( + reservations.get("waiting_waiters") or {}) + budget["fairness_waiter"] = copy.deepcopy( + reservations.get("fairness_waiter")) if snapshot is not None: effective_after = decision["effective_after"] budget["spendable_before_route"] = { diff --git a/promptpilot/worker.py b/promptpilot/worker.py index a91f43c..aad91f5 100644 --- a/promptpilot/worker.py +++ b/promptpilot/worker.py @@ -2310,9 +2310,21 @@ def _claim_next_task(busy_keys=(), busy_lane_ids=()): json.JSONDecodeError) as exc: print(f" !! pipeline lane scheduler unavailable: {exc}", flush=True) policy = None + try: + budget_fairness = pipeline_insights.worker_budget_fairness_policy() + except (AttributeError, OSError, TypeError, ValueError, + json.JSONDecodeError) as exc: + print(f" !! pipeline budget fairness unavailable: {exc}", flush=True) + budget_fairness = None + fairness_kwargs = ({ + "budget_wait_scope": budget_fairness["scope"], + "budget_starvation_timeout_seconds": ( + budget_fairness["starvation_timeout_seconds"]), + } if budget_fairness is not None else {}) if policy is None: return db.get_next_runnable( - busy_keys=busy_keys, key_fn=lock_key), None + busy_keys=busy_keys, key_fn=lock_key, + **fairness_kwargs), None assignments = {} @@ -2326,7 +2338,8 @@ def rank(task): return score task = db.get_next_runnable( - busy_keys=busy_keys, key_fn=lock_key, order_key_fn=rank) + busy_keys=busy_keys, key_fn=lock_key, order_key_fn=rank, + **fairness_kwargs) return task, assignments.get(task.id) if task is not None else None diff --git a/tests/test_github_budget.py b/tests/test_github_budget.py index bae940c..fad5396 100644 --- a/tests/test_github_budget.py +++ b/tests/test_github_budget.py @@ -14,7 +14,7 @@ import pytest -from promptpilot import api, bot, pipeline_insights, worker +from promptpilot import api, bot, db as database, pipeline_insights, worker from promptpilot.models import TaskCreate @@ -81,17 +81,48 @@ def _running_task(database, prompt): def _reserve(database, task, token, *, core=500, remaining=1000, - reset=None): + minimum=100, reset=None, starvation_timeout_seconds=0, + now=None): return database.reserve_pipeline_github_budget( "github-default", token=token, task_id=task.id, task_started_at=task.started_at, profile_id="example", queue_id="review", route="skill", cost={"core": core, "search": 0, "graphql": 0}, limits=_limits(core=remaining, reset=reset), - minimum_remaining={"core": 100, "search": 0, "graphql": 0}, + minimum_remaining={"core": minimum, "search": 0, "graphql": 0}, + starvation_timeout_seconds=starvation_timeout_seconds, + now=now, ) +def _budget_wait_row(database, task_id): + with database._connect() as conn: + return dict(conn.execute( + """SELECT budget_wait_scope, budget_wait_started_at + FROM tasks WHERE id = ?""", + (task_id,), + ).fetchone()) + + +def _arm_and_defer_budget_waiter( + database, task, *, wait_started_at, scope="github-default", + defer_until=None): + revision = database.pipeline_github_budget_reservations(scope)["revision"] + armed = database.arm_pipeline_github_budget_waiter( + scope, task_id=task.id, task_started_at=task.started_at, + expected_revision=revision, now=wait_started_at) + assert armed == { + "armed": True, "revision": revision, "revision_changed": False, + } + assert database.defer_task( + task.id, defer_until or ( + datetime.now(timezone.utc) - timedelta(seconds=1)), + "waiting for a GitHub reservation", budget_wait_scope=scope, + budget_wait_revision=revision, + expected_started_at=task.started_at) + return database.get_task(task.id) + + def test_budget_decision_defers_to_exact_latest_reset_plus_grace(): profile = _profile() policy = pipeline_insights._github_budget_policy(profile) @@ -198,6 +229,22 @@ def test_priority_one_headroom_defaults_to_zero_and_accepts_exact_vector(): } +def test_starvation_timeout_is_opt_in_and_strictly_bounded(): + assert pipeline_insights._github_budget_policy( + _profile())["starvation_timeout_seconds"] == 0 + configured = _profile() + configured["github_budget"]["starvation_timeout_seconds"] = 900 + assert pipeline_insights._github_budget_policy( + configured)["starvation_timeout_seconds"] == 900 + + for invalid in (True, -1, 86401, 1.5): + invalid_profile = _profile() + invalid_profile["github_budget"][ + "starvation_timeout_seconds"] = invalid + with pytest.raises(ValueError): + pipeline_insights._github_budget_policy(invalid_profile) + + def test_route_priority_promotion_is_attempt_fenced_internal_and_reversible( isolated_db): created = isolated_db.create_task(TaskCreate( @@ -345,6 +392,8 @@ def test_shared_github_scope_uses_strongest_profile_hard_reserve( high["github_budget"]["priority_one_headroom"] = { "core": 800, "search": 12, "graphql": 1000, } + low["github_budget"]["starvation_timeout_seconds"] = 900 + high["github_budget"]["starvation_timeout_seconds"] = 600 monkeypatch.setattr( pipeline_insights, "_profiles", lambda: {"low": low, "high": high}) @@ -358,6 +407,7 @@ def test_shared_github_scope_uses_strongest_profile_hard_reserve( assert policy["priority_one_headroom"] == { "core": 1200, "search": 12, "graphql": 1000, } + assert policy["starvation_timeout_seconds"] == 600 monkeypatch.setattr( pipeline_insights, "_github_rate_limits", lambda: _limits(core=1000, search=30, graphql=5000)) @@ -1902,10 +1952,440 @@ def test_explicit_reschedule_releases_budget_priority_handoff(isolated_db): with isolated_db._connect() as conn: stored = conn.execute( - "SELECT budget_wait_scope FROM tasks WHERE id = ?", + """SELECT budget_wait_scope, budget_wait_started_at + FROM tasks WHERE id = ?""", (waiter.id,), ).fetchone() assert stored["budget_wait_scope"] is None + assert stored["budget_wait_started_at"] is None + + +def test_budget_waiter_fairness_activates_at_exact_timeout(isolated_db): + created = isolated_db.create_task(TaskCreate( + prompt="Example - FIX", recurrence="4h", priority=3)) + waiter = isolated_db.get_next_runnable() + assert waiter.id == created.id + _arm_and_defer_budget_waiter( + isolated_db, waiter, wait_started_at=1000, + defer_until=datetime.now(timezone.utc) + timedelta(hours=1)) + + before = isolated_db.pipeline_github_budget_reservations( + "github-default", starvation_timeout_seconds=600, now=1599) + boundary = isolated_db.pipeline_github_budget_reservations( + "github-default", starvation_timeout_seconds=600, now=1600) + + assert "fairness_waiter" not in before + assert boundary["fairness_waiter"] == { + "task_id": waiter.id, + "priority": 3, + "wait_started_at": datetime.fromtimestamp( + 1000, timezone.utc).isoformat(), + "eligible_at": datetime.fromtimestamp( + 1600, timezone.utc).isoformat(), + "wait_seconds": 600, + } + urgent = isolated_db.create_task(TaskCreate( + prompt="Example - MERGE", recurrence="4h", priority=1)) + claimed = isolated_db.get_next_runnable( + budget_wait_scope="github-default", + budget_starvation_timeout_seconds=600, + budget_fairness_now=1600) + assert claimed.id == waiter.id + assert claimed.id != urgent.id + assert claimed.next_run_at is None + + +def test_worker_lane_claim_gives_aged_waiter_a_slot( + isolated_db, monkeypatch): + profile = _profile_with_costs(core=200) + profile["github_budget"]["starvation_timeout_seconds"] = 600 + profile["queues"] = [ + {"id": "merge", "series_contains": " - MERGE"}, + {"id": "plan", "series_contains": " - PLAN"}, + ] + profile["scheduler"] = {"lanes": [ + {"id": "integration", "queues": ["merge"]}, + {"id": "intake", "queues": ["plan"]}, + ]} + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + + waiter_created = isolated_db.create_task(TaskCreate( + prompt="Example - PLAN", recurrence="4h", priority=4)) + waiter = isolated_db.get_next_runnable() + assert waiter.id == waiter_created.id + _arm_and_defer_budget_waiter( + isolated_db, waiter, wait_started_at=time.time() - 601) + urgent = isolated_db.create_task(TaskCreate( + prompt="Example - MERGE", recurrence="4h", priority=1)) + + claimed, lane_id = worker._claim_next_task() + + assert claimed.id == waiter.id + assert claimed.id != urgent.id + assert lane_id == "example:intake" + + +def test_fairness_rejects_closed_lanes_and_falls_back_to_claim_renamed_waiter( + isolated_db, monkeypatch): + profile = _profile_with_costs(core=200) + profile["github_budget"]["starvation_timeout_seconds"] = 600 + profile["queues"] = [ + {"id": "merge", "series_contains": " - MERGE"}, + ] + profile["scheduler"] = {"lanes": [ + {"id": "integration", "queues": ["merge"], "borrow": False}, + ]} + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + + waiter_created = isolated_db.create_task(TaskCreate( + prompt="Removed queue - FIX", recurrence="4h", priority=4)) + waiter = isolated_db.get_next_runnable() + assert waiter.id == waiter_created.id + _arm_and_defer_budget_waiter( + isolated_db, waiter, wait_started_at=time.time() - 601) + urgent = isolated_db.create_task(TaskCreate( + prompt="Example - MERGE", recurrence="4h", priority=1)) + + with pytest.raises(ValueError, match="borrow:true recovery lane"): + pipeline_insights.worker_lane_policy() + claimed, lane_id = worker._claim_next_task() + + assert claimed.id == waiter.id + assert claimed.id != urgent.id + assert lane_id is None + + +def test_budget_waiter_fairness_is_oldest_first_then_advances( + isolated_db): + first_created = isolated_db.create_task(TaskCreate( + prompt="Example - PLAN", recurrence="4h", priority=4)) + first = isolated_db.get_next_runnable() + assert first.id == first_created.id + _arm_and_defer_budget_waiter( + isolated_db, first, wait_started_at=1000) + + second_created = isolated_db.create_task(TaskCreate( + prompt="Example - FIX", recurrence="4h", priority=3)) + second = isolated_db.get_next_runnable() + assert second.id == second_created.id + _arm_and_defer_budget_waiter( + isolated_db, second, wait_started_at=1001) + + ledger = isolated_db.pipeline_github_budget_reservations( + "github-default", starvation_timeout_seconds=600, now=2000) + assert ledger["fairness_waiter"]["task_id"] == first.id + assert ledger["fairness_waiter"]["priority"] == 4 + + assert isolated_db.update_task_fields( + first.id, {"scheduled_at": datetime.now(timezone.utc)}) + advanced = isolated_db.pipeline_github_budget_reservations( + "github-default", starvation_timeout_seconds=600, now=2000) + assert advanced["fairness_waiter"]["task_id"] == second.id + + +def test_starved_waiter_baton_blocks_early_and_atomic_new_reservations( + isolated_db): + owner = _running_task(isolated_db, "Existing GitHub owner") + assert _reserve( + isolated_db, owner, "existing-owner", core=500, + remaining=2000)["allowed"] is True + + waiter_created = isolated_db.create_task(TaskCreate( + prompt="Example - FIX", recurrence="4h", priority=3)) + waiter = isolated_db.get_next_runnable() + assert waiter.id == waiter_created.id + _arm_and_defer_budget_waiter( + isolated_db, waiter, wait_started_at=1000) + + urgent_created = isolated_db.create_task(TaskCreate( + prompt="Example - MERGE", recurrence="4h", priority=1)) + urgent = isolated_db.get_next_runnable() + assert urgent.id == urgent_created.id + ledger = isolated_db.pipeline_github_budget_reservations( + "github-default", starvation_timeout_seconds=600, now=1600) + profile = _profile_with_costs(core=200) + profile["github_budget"]["starvation_timeout_seconds"] = 600 + policy = pipeline_insights._github_budget_policy(profile) + + early = pipeline_insights._priority_waiter_decision( + policy, _limits(core=2000), ledger, priority=urgent.priority, + task_id=urgent.id, status_revision=7) + assert early["state"] == "fairness_waiter" + assert early["fairness_waiter"]["task_id"] == waiter.id + assert pipeline_insights._priority_waiter_decision( + policy, _limits(core=2000), ledger, priority=3, + task_id=waiter.id, status_revision=7) is None + + atomic = _reserve( + isolated_db, urgent, "urgent-after-timeout", core=200, + remaining=2000, minimum=100, starvation_timeout_seconds=600, + now=1600) + assert atomic["allowed"] is False + assert atomic["state"] == "fairness_waiter" + assert atomic["fairness_waiter"]["task_id"] == waiter.id + + assert isolated_db.mark_cancelled( + urgent.id, expected_started_at=urgent.started_at) + assert isolated_db.release_pipeline_github_budget( + "github-default", token="existing-owner", task_id=owner.id, + task_started_at=owner.started_at) + selected = isolated_db.get_next_runnable() + assert selected.id == waiter.id + admitted = _reserve( + isolated_db, selected, "selected-starved-waiter", core=500, + remaining=1300, minimum=700, starvation_timeout_seconds=600, + now=1601) + assert admitted["allowed"] is True + assert admitted["minimum_remaining"]["core"] == 700 + assert admitted["effective_after"]["core"] == 800 + assert _budget_wait_row(isolated_db, waiter.id) == { + "budget_wait_scope": None, "budget_wait_started_at": None, + } + + +def test_starved_baton_keeps_floor_and_low_wait_releases_baton(isolated_db): + created = isolated_db.create_task(TaskCreate( + prompt="Example - PLAN", recurrence="4h", priority=4)) + waiter = isolated_db.get_next_runnable() + assert waiter.id == created.id + _arm_and_defer_budget_waiter( + isolated_db, waiter, wait_started_at=1000) + claimed = isolated_db.get_next_runnable() + assert claimed.id == waiter.id + + denied = _reserve( + isolated_db, claimed, "starved-but-low", core=500, + remaining=1100, minimum=700, starvation_timeout_seconds=600, + now=1600) + assert denied["allowed"] is False + assert denied["state"] == "low" + assert denied["minimum_remaining"]["core"] == 700 + assert _budget_wait_row(isolated_db, waiter.id) == { + "budget_wait_scope": None, "budget_wait_started_at": None, + } + + assert isolated_db.defer_task( + claimed.id, datetime.now(timezone.utc) + timedelta(minutes=5), + "live GitHub quota is low", expected_started_at=claimed.started_at) + assert _budget_wait_row(isolated_db, waiter.id) == { + "budget_wait_scope": None, "budget_wait_started_at": None, + } + + +def test_budget_waiter_clear_is_exact_attempt_fenced(isolated_db): + task = _running_task(isolated_db, "Example - FIX") + revision = isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] + assert isolated_db.arm_pipeline_github_budget_waiter( + "github-default", task_id=task.id, task_started_at=task.started_at, + expected_revision=revision, now=1000)["armed"] + + assert not isolated_db.clear_pipeline_github_budget_waiter( + task_id=task.id, + task_started_at=task.started_at + timedelta(microseconds=1)) + assert _budget_wait_row( + isolated_db, task.id)["budget_wait_scope"] == "github-default" + assert isolated_db.clear_pipeline_github_budget_waiter( + task_id=task.id, + task_started_at=task.started_at) + assert _budget_wait_row(isolated_db, task.id) == { + "budget_wait_scope": None, "budget_wait_started_at": None, + } + + +@pytest.mark.parametrize(("with_costs", "retain_budget"), [ + (False, True), + (True, False), +]) +def test_success_without_reservation_clears_aged_waiter( + isolated_db, monkeypatch, with_costs, retain_budget): + profile = (_profile_with_costs(core=100) + if with_costs else _profile()) + profile["github_budget"]["starvation_timeout_seconds"] = 600 + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "_github_rate_limits", lambda: _limits()) + pipeline_insights.release_execution_admission() + task = _running_task(isolated_db, "Example - REVIEW") + revision = isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] + assert isolated_db.arm_pipeline_github_budget_waiter( + "github-default", task_id=task.id, task_started_at=task.started_at, + expected_revision=revision, now=time.time() - 601)["armed"] + + route = pipeline_insights.execution_route( + task, task.prompt, retain_budget=retain_budget) + + assert route["action"] == "prompt" + assert _budget_wait_row(isolated_db, task.id) == { + "budget_wait_scope": None, "budget_wait_started_at": None, + } + assert isolated_db.pipeline_github_budget_reservations( + "github-default")["count"] == 0 + pipeline_insights.release_execution_admission() + + +def test_unmatched_route_clears_stale_waiter_from_renamed_profile( + isolated_db, monkeypatch): + task = _running_task(isolated_db, "Renamed pipeline series") + revision = isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] + assert isolated_db.arm_pipeline_github_budget_waiter( + "github-default", task_id=task.id, task_started_at=task.started_at, + expected_revision=revision, now=time.time() - 601)["armed"] + monkeypatch.setattr(pipeline_insights, "_matching_queue", lambda _task: None) + + route = pipeline_insights.execution_route(task, task.prompt) + + assert route == { + "action": "prompt", "mode": "skill", "prompt": task.prompt, + } + assert _budget_wait_row(isolated_db, task.id) == { + "budget_wait_scope": None, "budget_wait_started_at": None, + } + + +@pytest.mark.parametrize("case", [ + "policy_none", "invalid_config", "acquire_error", "lease_unavailable", +]) +def test_early_non_handoff_paths_clear_stale_waiter( + isolated_db, monkeypatch, case): + profile = _profile() + profile["github_budget"]["starvation_timeout_seconds"] = 600 + if case == "policy_none": + profile.pop("github_budget") + elif case == "invalid_config": + profile["github_budget"]["starvation_timeout_seconds"] = True + elif case == "acquire_error": + monkeypatch.setattr( + isolated_db, "acquire_pipeline_scan_lease", + lambda *_args, **_kwargs: (_ for _ in ()).throw( + sqlite3.OperationalError("database is locked"))) + else: + monkeypatch.setattr( + isolated_db, "acquire_pipeline_scan_lease", + lambda *_args, **_kwargs: { + "acquired": False, "state": "unavailable", + "status_revision": 0, + }) + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "_github_rate_limits", lambda: _limits()) + task = _running_task(isolated_db, "Example - REVIEW") + revision = isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] + assert isolated_db.arm_pipeline_github_budget_waiter( + "github-default", task_id=task.id, task_started_at=task.started_at, + expected_revision=revision, now=time.time() - 601)["armed"] + + route = pipeline_insights.execution_route(task, task.prompt) + + assert route["action"] == ("prompt" if case == "policy_none" else "defer") + assert _budget_wait_row(isolated_db, task.id) == { + "budget_wait_scope": None, "budget_wait_started_at": None, + } + + +def test_budget_wait_timestamp_survives_rearm_and_crash_recovery_but_not_edit( + isolated_db): + task = _running_task(isolated_db, "Example - FIX") + revision = isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] + assert isolated_db.arm_pipeline_github_budget_waiter( + "github-default", task_id=task.id, task_started_at=task.started_at, + expected_revision=revision, now=1000)["armed"] + assert isolated_db.arm_pipeline_github_budget_waiter( + "github-default", task_id=task.id, task_started_at=task.started_at, + expected_revision=revision, now=1200)["armed"] + expected_timestamp = datetime.fromtimestamp( + 1000, timezone.utc).isoformat() + assert _budget_wait_row(isolated_db, task.id) == { + "budget_wait_scope": "github-default", + "budget_wait_started_at": expected_timestamp, + } + + assert isolated_db.recover_running_attempt(task.id, task.started_at) + assert _budget_wait_row(isolated_db, task.id) == { + "budget_wait_scope": "github-default", + "budget_wait_started_at": expected_timestamp, + } + edited_after = datetime.now(timezone.utc) + assert isolated_db.update_task_fields(task.id, {"priority": 4}) + edited = _budget_wait_row(isolated_db, task.id) + assert edited["budget_wait_scope"] == "github-default" + restarted_at = datetime.fromisoformat(edited["budget_wait_started_at"]) + assert edited_after <= restarted_at <= datetime.now(timezone.utc) + + +@pytest.mark.parametrize("editor", ["task_priority", "series_priority"]) +def test_all_priority_edit_paths_restart_budget_wait_age( + isolated_db, editor): + created = isolated_db.create_task(TaskCreate( + prompt="Example - FIX", recurrence="4h", priority=3)) + task = isolated_db.get_next_runnable() + assert task.id == created.id + _arm_and_defer_budget_waiter( + isolated_db, task, wait_started_at=1000, + defer_until=datetime.now(timezone.utc) + timedelta(minutes=5)) + edited_after = datetime.now(timezone.utc) + + if editor == "task_priority": + assert isolated_db.update_priority(task.id, 4) + else: + assert isolated_db.update_series(task.series_id, {"priority": 4}) + + edited = _budget_wait_row(isolated_db, task.id) + assert edited["budget_wait_scope"] == "github-default" + restarted_at = datetime.fromisoformat(edited["budget_wait_started_at"]) + assert edited_after <= restarted_at <= datetime.now(timezone.utc) + + +def test_budget_waiter_fairness_migrates_legacy_database( + tmp_path, monkeypatch): + db_path = tmp_path / "legacy-promptpilot.db" + monkeypatch.setattr(database, "DB_DIR", tmp_path) + monkeypatch.setattr(database, "DB_PATH", db_path) + legacy_schema = database.SCHEMA.replace( + " budget_wait_started_at TEXT,\n", "") + assert legacy_schema != database.SCHEMA + with sqlite3.connect(db_path) as conn: + conn.executescript(legacy_schema) + conn.execute( + """INSERT INTO tasks + (prompt, status, priority, created_at, budget_wait_scope) + VALUES (?, 'pending', 3, ?, ?)""", + ("legacy waiter", datetime.now(timezone.utc).isoformat(), + "github-default"), + ) + + database.init_db() + database.init_db() + + with sqlite3.connect(db_path) as conn: + columns = {row[1] for row in conn.execute( + "PRAGMA table_info(tasks)").fetchall()} + indexes = {row[1] for row in conn.execute( + "PRAGMA index_list(tasks)").fetchall()} + migrated = conn.execute( + """SELECT budget_wait_scope, budget_wait_started_at + FROM tasks WHERE prompt = 'legacy waiter'""").fetchone() + marker = conn.execute( + "SELECT 1 FROM schema_migrations WHERE version = ?", + (database.PIPELINE_BUDGET_WAITER_FAIRNESS_SCHEMA_VERSION,), + ).fetchone() + integrity = conn.execute("PRAGMA integrity_check").fetchone()[0] + + assert "budget_wait_started_at" in columns + assert "idx_tasks_budget_wait_fairness" in indexes + assert migrated[0] == "github-default" + migrated_at = datetime.fromisoformat(migrated[1]) + assert migrated_at.tzinfo is not None + assert marker == (1,) + assert integrity == "ok" def test_stale_prune_yields_to_woken_higher_priority_waiter(isolated_db):