diff --git a/README.md b/README.md index 17447ad..15ca910 100644 --- a/README.md +++ b/README.md @@ -1076,7 +1076,11 @@ skill остаётся дорогим: стартовые `3750 + 250` Core со При действительно низком живом лимите задача откладывается до reset ресурса плюс `reset_grace_seconds`. Если лимит достаточен, но временно обещан другой -задаче, повтор происходит через `busy_retry_seconds`, без часовой дыры. +задаче, ожидание остаётся прерываемым: освобождение точной резервации сразу +будит такие задачи, а `busy_retry_seconds` служит страховочным опросом после +аварии процесса. Приоритетная эстафета сохраняется до фактического резервирования +бюджета, поэтому уже запущенный менее приоритетный этап не обгоняет разбуженный. +Жёсткий reset живого лимита этим сигналом не сокращается. Недоступный лимит, повреждённый ledger или некорректная конфигурация работают fail-closed на `unavailable_retry_seconds`. Для профилей с одним GitHub token и `PP_DATA_DIR` PromptPilot покомпонентно применяет самый большой настроенный hard diff --git a/promptpilot/db.py b/promptpilot/db.py index f643b5f..9273793 100644 --- a/promptpilot/db.py +++ b/promptpilot/db.py @@ -71,6 +71,7 @@ worktree_path TEXT, worktree_branch TEXT, note TEXT, + budget_wait_scope TEXT, verdict TEXT ,series_id INTEGER REFERENCES task_series(id) ); @@ -390,6 +391,8 @@ "ALTER TABLE pipeline_target_reservations ADD COLUMN herdr_pane_id TEXT NOT NULL DEFAULT ''", "ALTER TABLE pipeline_target_reservations ADD COLUMN herdr_tab_id TEXT NOT NULL DEFAULT ''", "ALTER TABLE pipeline_target_reservations ADD COLUMN herdr_workspace_id TEXT NOT NULL DEFAULT ''", + "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)", ] WORKFLOW_SCHEMA_VERSION = "workflow_orchestrator_w0_v1" @@ -425,6 +428,10 @@ def _parse_dt(val: Optional[str]) -> Optional[datetime]: def _row_to_task(row: sqlite3.Row) -> TaskInDB: d = dict(row) + # Internal scheduler metadata is deliberately not part of the public task + # model/API. It only identifies interruptible waits on another task's + # GitHub budget reservation. + d.pop("budget_wait_scope", None) for field in ("scheduled_at", "next_run_at", "created_at", "started_at", "completed_at"): d[field] = _parse_dt(d[field]) return TaskInDB(**d) @@ -1267,7 +1274,8 @@ def mark_completed(task_id: int, result: str, exit_code: int = 0, "UPDATE tasks SET status = 'completed', result = ?, error = NULL, " "next_run_at = NULL, exit_code = ?, completed_at = ?, model_used = ?, " "session_id = COALESCE(?, session_id), " - f"verdict = COALESCE(?, verdict), note = NULL WHERE {where}", + f"verdict = COALESCE(?, verdict), note = NULL, " + f"budget_wait_scope = NULL WHERE {where}", values, ) if cur.rowcount: @@ -1287,7 +1295,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 WHERE {where}", + f"completed_at = ?, note = NULL, budget_wait_scope = NULL " + f"WHERE {where}", values, ) if cur.rowcount: @@ -1312,7 +1321,7 @@ 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 + completed_at = ?, note = NULL, budget_wait_scope = NULL WHERE id = ? AND status = 'running' AND started_at = ?""", (error, exit_code, _now(), task_id, started_at), ) @@ -1336,7 +1345,7 @@ 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) + error = COALESCE(?, error), budget_wait_scope = NULL WHERE """ + where, values, ) @@ -1348,6 +1357,8 @@ def mark_rate_limited(task_id: int, next_run_at: datetime, error: str = None, def defer_task(task_id: int, next_run_at: datetime, reason: str = None, *, hard_not_before: bool = False, + budget_wait_scope: str | None = None, + budget_wait_revision: int | None = None, expected_started_at=None) -> bool: """Return a claimed task to pending without consuming a retry attempt. @@ -1356,19 +1367,48 @@ def defer_task(task_id: int, next_run_at: datetime, reason: str = None, move. A hard deadline is additionally stored in ``next_run_at`` so automatic queue wake-ups cannot spend GitHub quota before the exact reset; an explicit human ``run_now`` still clears that barrier intentionally. + + ``budget_wait_scope`` marks the distinct, interruptible case where live + quota is sufficient but another task temporarily owns a reservation. The + ledger release wakes that waiter immediately instead of leaving a + high-priority task asleep for the whole busy-retry interval. """ + if budget_wait_scope is not None: + if (not isinstance(budget_wait_scope, str) + or not budget_wait_scope.strip() + or len(budget_wait_scope) > 256): + raise ValueError("budget wait scope must contain 1..256 characters") + if (isinstance(budget_wait_revision, bool) + or not isinstance(budget_wait_revision, int) + or budget_wait_revision < 0): + raise ValueError( + "budget wait revision must be a non-negative integer") + if hard_not_before: + raise ValueError( + "budget reservation wait cannot also be a hard deadline") + elif budget_wait_revision is not None: + raise ValueError("budget wait revision requires a scope") deadline = _to_utc_iso(next_run_at) - with _connect() as conn: + with _connect(immediate=budget_wait_scope is not None) as conn: + # Close the release-before-defer race. If the reservation ledger moved + # after admission observed contention, the sole wake may already have + # happened. Requeue immediately, but preserve the durable priority + # handoff until this task actually reserves budget. + if (budget_wait_scope is not None + and _pipeline_github_budget_revision( + conn, budget_wait_scope) != budget_wait_revision): + deadline = _now() attempt = _attempt_iso(expected_started_at) where = "id = ? AND status = 'running'" - values = [deadline, deadline if hard_not_before else None, reason, task_id] + values = [deadline, deadline if hard_not_before else None, reason, + budget_wait_scope, task_id] if attempt is not None: where += " AND started_at = ?" values.append(attempt) cur = conn.execute( """UPDATE tasks SET status = 'pending', scheduled_at = ?, next_run_at = ?, started_at = NULL, - error = COALESCE(?, error) + error = COALESCE(?, error), budget_wait_scope = ? WHERE """ + where, values, ) @@ -1413,7 +1453,8 @@ def mark_cancelled(task_id: int, note: str = None, values.append(attempt) cur = conn.execute( "UPDATE tasks SET status = 'cancelled', completed_at = ?, " - f"error = COALESCE(?, error), note = NULL WHERE {where}", + f"error = COALESCE(?, error), note = NULL, " + f"budget_wait_scope = NULL WHERE {where}", values, ) if cur.rowcount: @@ -1425,7 +1466,9 @@ def mark_cancelled(task_id: int, note: str = None, def cancel_task(task_id: int) -> bool: with _connect() as conn: cur = conn.execute( - "UPDATE tasks SET status = 'cancelled', completed_at = ? WHERE id = ? AND status IN ('pending', 'rate_limited')", + "UPDATE tasks SET status = 'cancelled', completed_at = ?, " + "budget_wait_scope = NULL WHERE id = ? " + "AND status IN ('pending', 'rate_limited')", (_now(), task_id), ) if cur.rowcount: @@ -2018,7 +2061,8 @@ def series_action(series_id: int, action: str) -> bool: elif action == "end": 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 = ? " + conn.execute("UPDATE tasks SET status = 'cancelled', completed_at = ?, " + "budget_wait_scope = NULL " "WHERE series_id = ? AND status IN ('pending', 'rate_limited')", (now, series_id)) conn.execute( @@ -2299,6 +2343,10 @@ def update_task_fields(task_id: int, fields: dict) -> bool: return False if "scheduled_at" in fields: 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 sets = ", ".join(f"{k} = ?" for k in fields) with _connect() as conn: cur = conn.execute( @@ -2709,6 +2757,8 @@ def _pipeline_scan_lease_key(scope: str) -> str: _PIPELINE_GITHUB_BUDGET_RESERVATION_PREFIX = \ "pipeline_github_budget_reservations:v1:" +_PIPELINE_GITHUB_BUDGET_REVISION_PREFIX = \ + "pipeline_github_budget_revision:v1:" _PIPELINE_GITHUB_RATE_SNAPSHOT_PREFIX = "pipeline_github_rate_snapshot:v1:" _PIPELINE_GITHUB_BUDGET_RESOURCES = ("core", "search", "graphql") @@ -2721,6 +2771,40 @@ def _pipeline_github_budget_reservation_key(scope: str) -> str: return f"{_PIPELINE_GITHUB_BUDGET_RESERVATION_PREFIX}{digest}" +def _pipeline_github_budget_revision_key(scope: str) -> str: + value = str(scope).strip() + if not value: + raise ValueError("pipeline GitHub budget scope must not be empty") + digest = hashlib.sha256(value.encode("utf-8")).hexdigest() + return f"{_PIPELINE_GITHUB_BUDGET_REVISION_PREFIX}{digest}" + + +def _pipeline_github_budget_revision(conn, scope: str) -> int: + key = _pipeline_github_budget_revision_key(scope) + row = conn.execute( + "SELECT value FROM settings WHERE key = ?", (key,) + ).fetchone() + if row is None: + return 0 + try: + revision = int(row["value"]) + except (TypeError, ValueError, OverflowError) as exc: + raise ValueError( + "pipeline GitHub budget revision is corrupt") from exc + if revision < 0: + raise ValueError("pipeline GitHub budget revision is corrupt") + return revision + + +def _advance_pipeline_github_budget_revision(conn, scope: str) -> int: + revision = _pipeline_github_budget_revision(conn, scope) + 1 + conn.execute( + "INSERT OR REPLACE INTO settings (key, value) VALUES (?, ?)", + (_pipeline_github_budget_revision_key(scope), str(revision)), + ) + return revision + + def _pipeline_github_rate_snapshot_key(scope: str) -> str: value = str(scope).strip() if not value: @@ -2802,10 +2886,11 @@ def _load_pipeline_budget_reservations(conn, key: str) -> list[dict]: raise ValueError("pipeline GitHub budget reservation ledger is corrupt") from exc -def _load_or_recover_pipeline_budget_reservations(conn, key: str) -> list[dict]: +def _load_or_recover_pipeline_budget_reservations( + conn, key: str, scope: str) -> tuple[list[dict], dict | None]: """Reclaim corrupt state only when no provider attempt can still own it.""" try: - return _load_pipeline_budget_reservations(conn, key) + return _load_pipeline_budget_reservations(conn, key), None except ValueError: running = conn.execute( "SELECT 1 FROM tasks WHERE status = 'running' LIMIT 1" @@ -2813,7 +2898,9 @@ def _load_or_recover_pipeline_budget_reservations(conn, key: str) -> list[dict]: if running is not None: raise conn.execute("DELETE FROM settings WHERE key = ?", (key,)) - return [] + _advance_pipeline_github_budget_revision(conn, scope) + woken = _wake_pipeline_github_budget_waiters(conn, scope) + return [], woken def _active_pipeline_budget_reservations(conn, values: list[dict]) -> tuple[list[dict], bool]: @@ -2841,19 +2928,110 @@ def _write_pipeline_budget_reservations(conn, key: str, values: list[dict]) -> N ) +def _pipeline_github_budget_waiter_summary(conn, scope: str) -> dict: + row = conn.execute( + """SELECT COUNT(*) AS count, MIN(priority) AS min_priority + FROM tasks + WHERE budget_wait_scope = ? 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))))""", + (scope,), + ).fetchone() + 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, + } + + +def _wake_pipeline_github_budget_waiters(conn, scope: str) -> dict: + """Wake soft waits while preserving their priority handoff marker.""" + summary = _pipeline_github_budget_waiter_summary(conn, scope) + cur = conn.execute( + """UPDATE tasks + SET scheduled_at = ?, next_run_at = NULL + WHERE status = 'pending' AND budget_wait_scope = ?""", + (_now(), scope), + ) + return { + "count": cur.rowcount, + "min_priority": summary["min_priority"], + } + + +def arm_pipeline_github_budget_waiter( + scope: str, *, task_id: int, task_started_at, + expected_revision: int | None = None) -> dict: + """Persist a priority handoff before releasing the admission scan lease. + + New reservations are serialized by that lease, while reservation release + is intentionally independent. Recording the exact running attempt before + the lease is released closes both release-before-defer and + release-before-lower-admission windows. The marker survives wake and claim; + it is cleared only by successful reservation, terminal settlement, a + different defer reason, or an explicit operator reschedule. + """ + if not isinstance(scope, str) or not scope.strip() or len(scope) > 256: + raise ValueError("pipeline GitHub budget scope must contain 1..256 characters") + if (expected_revision is not None + and (isinstance(expected_revision, bool) + or not isinstance(expected_revision, int) + or expected_revision < 0)): + raise ValueError( + "expected pipeline budget revision must be a non-negative integer") + attempt = _attempt_iso(task_started_at) + if attempt is None: + raise ValueError("running task attempt is required for budget waiter") + with _connect(immediate=True) as conn: + revision = _pipeline_github_budget_revision(conn, scope) + cur = conn.execute( + """UPDATE tasks SET budget_wait_scope = ? + WHERE id = ? AND status = 'running' AND started_at = ?""", + (scope, task_id, attempt), + ) + return { + "armed": cur.rowcount > 0, + "revision": revision, + "revision_changed": ( + expected_revision is not None + and revision != expected_revision), + } + + def pipeline_github_budget_reservations(scope: str) -> dict: """Return live, task-fenced reservations and prune completed attempts.""" key = _pipeline_github_budget_reservation_key(scope) with _connect(immediate=True) as conn: - values = _load_or_recover_pipeline_budget_reservations(conn, key) + values, recovered = _load_or_recover_pipeline_budget_reservations( + conn, key, scope) active, changed = _active_pipeline_budget_reservations(conn, values) + revision = _pipeline_github_budget_revision(conn, scope) + woken = recovered or {"count": 0, "min_priority": None} if changed: _write_pipeline_budget_reservations(conn, key, active) + revision = _advance_pipeline_github_budget_revision(conn, scope) + stale_woken = _wake_pipeline_github_budget_waiters(conn, scope) + priorities = [ + value for value in ( + woken.get("min_priority"), + stale_woken.get("min_priority")) + if type(value) is int + ] + woken = { + "count": int(woken.get("count") or 0) + + int(stale_woken.get("count") or 0), + "min_priority": min(priorities) if priorities else None, + } totals = {resource: 0 for resource in _PIPELINE_GITHUB_BUDGET_RESOURCES} for item in active: for resource in totals: totals[resource] += item["cost"][resource] - return {"count": len(active), "totals": totals, "items": active} + waiting = _pipeline_github_budget_waiter_summary(conn, scope) + return {"count": len(active), "totals": totals, "items": active, + "revision": revision, "woken_waiters": woken, + "waiting_waiters": waiting} def record_pipeline_github_rate_snapshot( @@ -2945,7 +3123,8 @@ def reserve_pipeline_github_budget( profile_id: str, queue_id: str, route: str, cost: dict, limits: dict, minimum_remaining: dict, now: Optional[float] = None, - scan_lease_guard: Optional[dict] = None) -> dict: + scan_lease_guard: Optional[dict] = None, + expected_revision: Optional[int] = None) -> dict: """Atomically reserve quota for one exact running task attempt.""" key = _pipeline_github_budget_reservation_key(scope) if not _valid_pipeline_scan_token(token): @@ -2958,6 +3137,12 @@ def reserve_pipeline_github_budget( or not isinstance(queue_id, str) or not queue_id or not isinstance(route, str) or not route): raise ValueError("invalid pipeline budget reservation identity") + if (expected_revision is not None + and (isinstance(expected_revision, bool) + or not isinstance(expected_revision, int) + or expected_revision < 0)): + raise ValueError( + "expected pipeline budget revision must be a non-negative integer") requested = _pipeline_budget_vector(cost, "cost") floor = _pipeline_budget_vector(minimum_remaining, "minimum_remaining") reported = {} @@ -2991,14 +3176,33 @@ def reserve_pipeline_github_budget( return {"allowed": False, "state": "scan_lease_lost", "reason": "GitHub scan lease changed before budget reservation"} task = conn.execute( - "SELECT status, started_at FROM tasks WHERE id = ?", (task_id,) + "SELECT status, started_at, priority FROM tasks WHERE id = ?", + (task_id,) ).fetchone() if (task is None or task["status"] != "running" or task["started_at"] != task_started_at): return {"allowed": False, "state": "task_fence_lost", "reason": "running task attempt changed before budget reservation"} - values = _load_or_recover_pipeline_budget_reservations(conn, key) + values, recovered = _load_or_recover_pipeline_budget_reservations( + conn, key, scope) active, changed = _active_pipeline_budget_reservations(conn, values) + revision = _pipeline_github_budget_revision(conn, scope) + woken = recovered or {"count": 0, "min_priority": None} + if changed: + _write_pipeline_budget_reservations(conn, key, active) + revision = _advance_pipeline_github_budget_revision(conn, scope) + stale_woken = _wake_pipeline_github_budget_waiters(conn, scope) + priorities = [ + value for value in ( + woken.get("min_priority"), + stale_woken.get("min_priority")) + if type(value) is int + ] + woken = { + "count": int(woken.get("count") or 0) + + int(stale_woken.get("count") or 0), + "min_priority": min(priorities) if priorities else None, + } own = [item for item in active if item["token"] == token] if own: item = own[0] @@ -3006,8 +3210,12 @@ def reserve_pipeline_github_budget( or item["task_started_at"] != task_started_at or item["cost"] != requested or item["route"] != route): raise ValueError("pipeline budget token was reused inconsistently") - if changed: - _write_pipeline_budget_reservations(conn, key, active) + conn.execute( + """UPDATE tasks SET budget_wait_scope = NULL + WHERE id = ? AND status = 'running' AND started_at = ? + AND budget_wait_scope = ?""", + (task_id, task_started_at, scope), + ) return {"allowed": True, "state": "reserved", "reported_remaining": reported, "reserved_other": {resource: sum( @@ -3019,7 +3227,8 @@ def reserve_pipeline_github_budget( value["cost"][resource] for value in active) for resource in reported}, "blocked_resources": [], - "active_reservations": len(active)} + "active_reservations": len(active), + "revision": revision} if any(item["task_id"] == task_id and item["task_started_at"] == task_started_at for item in active): raise ValueError("running task attempt already owns another reservation") @@ -3028,6 +3237,57 @@ def reserve_pipeline_github_budget( item["cost"][resource] for item in active) for resource in reported} after = {resource: reported[resource] - reserved[resource] - requested[resource] for resource in reported} + if expected_revision is not None and revision != expected_revision: + return { + "allowed": False, "state": "ledger_changed", + "reason": ( + "GitHub budget reservations changed between admission " + "and final reservation"), + "blocked_resources": [], + "reported_remaining": reported, "reserved_other": reserved, + "requested_cost": requested, "minimum_remaining": floor, + "effective_after": after, "active_reservations": len(active), + "revision": revision, + } + + waiting = _pipeline_github_budget_waiter_summary(conn, scope) + waiting_priority = waiting.get("min_priority") + if (type(waiting_priority) is int + and waiting_priority < int(task["priority"])): + return { + "allowed": False, "state": "priority_waiter", + "reason": ( + "a higher-priority task is waiting for this GitHub " + "budget reservation"), + "blocked_resources": [], + "reported_remaining": reported, "reserved_other": reserved, + "requested_cost": requested, "minimum_remaining": floor, + "effective_after": after, "active_reservations": len(active), + "revision": revision, + } + + # Crash recovery may discover and prune a stale owner while a lower- + # priority task is trying to reserve. Yield once to a waiter that the + # same transaction just woke, otherwise the lower-priority task can + # immediately consume the released budget and recreate the inversion. + woken_priority = woken.get("min_priority") + if ((changed or recovered is not None) + and type(woken_priority) is int + and woken_priority < int(task["priority"])): + return { + "allowed": False, "state": "priority_waiter", + "reason": ( + "a higher-priority GitHub budget waiter was woken after " + "stale reservation recovery"), + "blocked_resources": [], + "reported_remaining": reported, + "reserved_other": {resource: sum( + item["cost"][resource] for item in active) + for resource in reported}, + "requested_cost": requested, "minimum_remaining": floor, + "effective_after": {}, "active_reservations": len(active), + "revision": revision, + } blocked = [{ "resource": resource, "reported_remaining": reported[resource], @@ -3041,8 +3301,6 @@ def reserve_pipeline_github_budget( < floor[resource] else "reservation"), } for resource in reported if after[resource] < floor[resource]] if blocked: - if changed: - _write_pipeline_budget_reservations(conn, key, active) state = ("low" if any(item["blocked_by"] == "live" for item in blocked) else "budget_in_flight") @@ -3050,7 +3308,8 @@ def reserve_pipeline_github_budget( "blocked_resources": blocked, "reported_remaining": reported, "reserved_other": reserved, "requested_cost": requested, "minimum_remaining": floor, - "effective_after": after, "active_reservations": len(active)} + "effective_after": after, "active_reservations": len(active), + "revision": revision} active.append({ "token": token, "task_id": task_id, @@ -3059,11 +3318,19 @@ def reserve_pipeline_github_budget( "created_at": created_at, }) _write_pipeline_budget_reservations(conn, key, active) + revision = _advance_pipeline_github_budget_revision(conn, scope) + conn.execute( + """UPDATE tasks SET budget_wait_scope = NULL + WHERE id = ? AND status = 'running' AND started_at = ? + AND budget_wait_scope = ?""", + (task_id, task_started_at, scope), + ) return {"allowed": True, "state": "reserved", "reported_remaining": reported, "reserved_other": reserved, "requested_cost": requested, "minimum_remaining": floor, "effective_after": after, - "active_reservations": len(active)} + "active_reservations": len(active), + "revision": revision} def release_pipeline_github_budget( @@ -3075,7 +3342,8 @@ def release_pipeline_github_budget( if isinstance(task_started_at, datetime): task_started_at = _to_utc_iso(task_started_at) with _connect(immediate=True) as conn: - values = _load_or_recover_pipeline_budget_reservations(conn, key) + values, _recovered = _load_or_recover_pipeline_budget_reservations( + conn, key, scope) kept = [item for item in values if not ( item["token"] == token and item["task_id"] == task_id and item["task_started_at"] == task_started_at)] @@ -3084,6 +3352,8 @@ def release_pipeline_github_budget( if not removed and not changed: return False _write_pipeline_budget_reservations(conn, key, active) + _advance_pipeline_github_budget_revision(conn, scope) + _wake_pipeline_github_budget_waiters(conn, scope) return removed @@ -3651,7 +3921,8 @@ def recover_running_attempt(task_id: int, started_at) -> bool: cur = conn.execute( """UPDATE tasks SET status = 'cancelled', completed_at = ?, - error = COALESCE(error, ?), note = NULL + error = COALESCE(error, ?), note = NULL, + budget_wait_scope = 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 218395e..56993d1 100644 --- a/promptpilot/pipeline_insights.py +++ b/promptpilot/pipeline_insights.py @@ -696,6 +696,77 @@ def _evaluate_github_budget(policy: dict, limits: dict | None, *, } +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.""" + woken = reservations.get("woken_waiters") or {} + waiting = reservations.get("waiting_waiters") or {} + priorities = [ + value for value in ( + woken.get("min_priority"), waiting.get("min_priority")) + if type(value) is int + ] + waiter_priority = min(priorities) if priorities else None + if (type(priority) is not int or not 1 <= priority <= 10 + or type(waiter_priority) is not int + or waiter_priority >= priority): + return None + now = time.time() + decision = _budget_denied( + policy, state="priority_waiter", + reason=( + "GitHub API-бюджет освобождён; сначала разбужена задача " + f"с более высоким приоритетом {waiter_priority}"), + 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), + ) + revision = reservations.get("revision") + if type(revision) is int and revision >= 0: + decision["_budget_reservation_revision"] = revision + return decision + + +def _arm_budget_waiter(task, decision: dict) -> dict: + """Persist a budget handoff before the serial admission lease is released.""" + state = decision.get("state") + if state not in { + "budget_in_flight", "ledger_changed", "priority_waiter", + "scan_in_progress"}: + return decision + scope = decision.get("lease_scope") + revision = decision.get("_budget_reservation_revision") + task_id = getattr(task, "id", None) + started_at = getattr(task, "started_at", None) + revision_valid = type(revision) is int and revision >= 0 + if (not isinstance(scope, str) or not scope + or (state != "scan_in_progress" and not revision_valid) + or type(task_id) is not int or task_id <= 0 + or started_at is None): + raise ValueError( + "running task identity or reservation revision is unavailable " + "for GitHub budget wait") + armed = db.arm_pipeline_github_budget_waiter( + scope, task_id=task_id, task_started_at=started_at, + expected_revision=(revision if revision_valid else None)) + if not armed.get("armed"): + raise ValueError( + "running task attempt changed before GitHub budget wait was armed") + current_revision = armed.get("revision") + if type(current_revision) is not int or current_revision < 0: + raise ValueError("GitHub budget waiter returned an invalid revision") + decision["_budget_reservation_revision"] = current_revision + if armed.get("revision_changed"): + # A release may race the budget observation, but cannot race a lower + # admission while this scan lease is held. Keep the priority marker and + # make this exact task runnable immediately after it is deferred. + decision["defer_until"] = _defer_at(time.time()) + return decision + + @contextmanager def _github_scan_admission(profile: dict, purpose: str, *, profile_id: str | None = None, @@ -751,11 +822,23 @@ def _github_scan_admission(profile: dict, purpose: str, retry_at = now + policy["unavailable_retry_seconds"] reason = "SQLite lease GitHub-сканирования повреждена; scan запрещён" blocked_state = "lease_unavailable" - yield _budget_denied( + decision = _budget_denied( policy, state=blocked_state, reason=reason, now=now, defer_at=retry_at, status_revision=(status_revision if isinstance(status_revision, int) else None)) + if task is not None: + try: + decision = _arm_budget_waiter(task, decision) + except (TypeError, ValueError, sqlite3.Error) as exc: + decision = _budget_denied( + policy, state="lease_unavailable", + reason=f"GitHub budget handoff недоступен: {exc}", + now=now, defer_at=now + policy["unavailable_retry_seconds"], + status_revision=( + status_revision + if isinstance(status_revision, int) else None)) + yield decision return status_revision = acquired.get("status_revision") @@ -795,11 +878,23 @@ def _github_scan_admission(profile: dict, purpose: str, now=time.time(), limits=limits, status_revision=lease.status_revision) else: - decision = _evaluate_github_budget( - policy, limits, now=time.time(), - status_revision=lease.status_revision, - reserved_other=reservations["totals"], - active_reservations=reservations["count"]) + admission_priority = getattr( + task, "priority", + 10 if budget_route == "insights" else None) + decision = _priority_waiter_decision( + policy, limits, reservations, + priority=admission_priority, + status_revision=lease.status_revision) + if decision is None: + decision = _evaluate_github_budget( + policy, limits, now=time.time(), + status_revision=lease.status_revision, + reserved_other=reservations["totals"], + active_reservations=reservations["count"]) + decision["_budget_reservation_revision"] = \ + reservations["revision"] + if task is not None: + decision = _arm_budget_waiter(task, decision) except _GitHubScanLeaseFailure as exc: decision = _lease_exception_decision( profile, exc, status_revision=lease.status_revision) @@ -862,7 +957,17 @@ def _reservation_denied(policy: dict, limits: dict | None, result: dict, *, now = time.time() blocked = result.get("blocked_resources") or [] active_reservations = int(result.get("active_reservations") or 0) - if blocked: + if result.get("state") == "ledger_changed": + defer_at = now + 1 + state = "ledger_changed" + reason = str(result.get("reason") or + "GitHub budget changed during pipeline admission") + elif result.get("state") == "priority_waiter": + defer_at = now + policy["busy_retry_seconds"] + state = "priority_waiter" + reason = str(result.get("reason") or + "GitHub budget yielded to a higher-priority waiter") + elif blocked: live_blocked = [item for item in blocked if item.get("blocked_by") == "live"] if live_blocked: @@ -885,13 +990,17 @@ def _reservation_denied(policy: dict, limits: dict | None, result: dict, *, state = str(result.get("state") or "reservation_unavailable") reason = str(result.get("reason") or "GitHub API reservation не создана") - return _budget_denied( + decision = _budget_denied( policy, state=state, reason=reason, now=now, limits=limits, defer_at=defer_at, blocked_resources=blocked, status_revision=status_revision, reserved_other=result.get("reserved_other"), effective_after=result.get("effective_after"), active_reservations=active_reservations) + revision = result.get("revision") + if type(revision) is int and revision >= 0: + decision["_budget_reservation_revision"] = revision + return decision def _reserve_execution_admission( @@ -899,6 +1008,9 @@ def _reserve_execution_admission( budget_route: str, *, retain_budget: bool) -> dict: """Recheck the elected route and optionally reserve its in-flight cost.""" lease = admission.get("_lease") + expected_revision = admission.get("_budget_reservation_revision") + if type(expected_revision) is not int or expected_revision < 0: + expected_revision = None try: policy = _github_budget_policy(profile) if policy is None: @@ -933,11 +1045,18 @@ def _reserve_execution_admission( elif not retain_budget: reservations = db.pipeline_github_budget_reservations( policy["lease_scope"]) - decision = _evaluate_github_budget( - policy, limits, now=time.time(), - status_revision=lease.status_revision, - reserved_other=reservations["totals"], - active_reservations=reservations["count"]) + decision = _priority_waiter_decision( + policy, limits, reservations, + priority=getattr(task, "priority", None), + status_revision=lease.status_revision) + if decision is None: + decision = _evaluate_github_budget( + policy, limits, now=time.time(), + status_revision=lease.status_revision, + reserved_other=reservations["totals"], + active_reservations=reservations["count"]) + decision["_budget_reservation_revision"] = \ + reservations["revision"] else: task_id = getattr(task, "id", None) started_at = getattr(task, "started_at", None) @@ -951,7 +1070,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"], - scan_lease_guard=lease.guard) + scan_lease_guard=lease.guard, + expected_revision=expected_revision) if not reserved.get("allowed"): decision = _reservation_denied( policy, limits, reserved, @@ -963,6 +1083,8 @@ def _reserve_execution_admission( reserved_other=reserved.get("reserved_other"), active_reservations=int( reserved.get("active_reservations") or 0)) + decision["_budget_reservation_revision"] = \ + reserved.get("revision") if not decision.get("allowed"): db.release_pipeline_github_budget( policy["lease_scope"], token=token, task_id=task_id, @@ -979,6 +1101,7 @@ def _reserve_execution_admission( _execution_budget_context.reservation = \ _GitHubBudgetReservation( policy["lease_scope"], token, task_id, started_at) + decision = _arm_budget_waiter(task, decision) except _GitHubScanLeaseFailure as exc: decision = _lease_exception_decision( profile, exc, @@ -2513,15 +2636,34 @@ def _budget_defer_route(admission: dict, profile_id: str, profile: dict, if phase: context = _budget_defer_context(admission, phase, preflight) reason = f"{reason}; preflight: {_format_budget_defer_context(context)}" + reservation_revision = admission.get("_budget_reservation_revision") + reservation_handoff = ( + admission.get("state") in { + "budget_in_flight", "ledger_changed", "priority_waiter", + "scan_in_progress"} + and type(reservation_revision) is int + and reservation_revision >= 0 + and isinstance(admission.get("lease_scope"), str) + and bool(admission.get("lease_scope"))) + defer_policy = ( + "reservation_release" if ( + reservation_handoff + and admission.get("state") == "budget_in_flight") + else "retry" if admission.get("state") in { + "ledger_changed", "priority_waiter", "scan_in_progress"} + else "hard_not_before") result = { "action": "defer", "mode": mode, "reason": reason, "defer_until": admission.get("defer_until"), - "defer_policy": "hard_not_before", + "defer_policy": defer_policy, "profile_id": profile_id, "queue_id": queue.get("id"), "github_budget": _public_budget_decision(admission), "github_rate_limit": admission.get("github_rate_limit"), } + if reservation_handoff: + result["budget_wait_scope"] = admission.get("lease_scope") + result["budget_wait_revision"] = reservation_revision if context is not None: result["defer_context"] = context return result diff --git a/promptpilot/worker.py b/promptpilot/worker.py index e9505bc..a91f43c 100644 --- a/promptpilot/worker.py +++ b/promptpilot/worker.py @@ -1414,6 +1414,8 @@ def _execute_task_body(task, admission_complete=None): task, next_run, reason, hard_not_before=( route.get("defer_policy") == "hard_not_before"), + budget_wait_scope=route.get("budget_wait_scope"), + budget_wait_revision=route.get("budget_wait_revision"), ), f"отложить pipeline-задачу #{task.id}", ) diff --git a/tests/test_github_budget.py b/tests/test_github_budget.py index 1727638..86329a7 100644 --- a/tests/test_github_budget.py +++ b/tests/test_github_budget.py @@ -574,6 +574,9 @@ def test_legacy_budget_without_costs_keeps_floor_check_but_never_reserves( "count": 0, "totals": {"core": 0, "search": 0, "graphql": 0}, "items": [], + "revision": 0, + "woken_waiters": {"count": 0, "min_priority": None}, + "waiting_waiters": {"count": 0, "min_priority": None}, } @@ -719,6 +722,9 @@ def test_reservation_contention_retries_quickly_but_live_low_waits_for_reset( assert contention["action"] == "defer" assert contention["github_budget"]["state"] == "budget_in_flight" assert contention["github_budget"]["active_reservations"] == 1 + assert contention["defer_policy"] == "reservation_release" + assert contention["budget_wait_scope"] == "github-default" + assert type(contention["budget_wait_revision"]) is int assert timedelta(0) < contention_until - before <= timedelta(seconds=10) assert isolated_db.release_pipeline_github_budget( @@ -730,6 +736,8 @@ def test_reservation_contention_retries_quickly_but_live_low_waits_for_reset( assert genuinely_low["action"] == "defer" assert genuinely_low["github_budget"]["state"] == "low" + assert genuinely_low["defer_policy"] == "hard_not_before" + assert "budget_wait_scope" not in genuinely_low assert genuinely_low["defer_until"] == datetime.fromtimestamp( reset + 17, timezone.utc).isoformat() assert isolated_db.pipeline_github_budget_reservations( @@ -1545,6 +1553,664 @@ def test_worker_uses_exact_budget_defer_without_loading_provider( assert deferred.error == reason +def test_budget_reservation_release_wakes_deferred_worker_immediately( + isolated_db, monkeypatch): + owner = _running_task(isolated_db, "Budget owner") + assert _reserve( + isolated_db, owner, "wake-owner", core=500, + remaining=1000)["allowed"] is True + created = isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + assert waiter is not None and waiter.id == created.id + revision = isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] + target = datetime.now(timezone.utc) + timedelta(minutes=7) + reason = "GitHub API budget is temporarily reserved by another task" + monkeypatch.setattr(pipeline_insights, "dispatch_gate", lambda _task: None) + monkeypatch.setattr( + pipeline_insights, "execution_route", lambda *_args, **_kwargs: { + "action": "defer", "mode": "skill", "reason": reason, + "defer_until": target.isoformat(), + "defer_policy": "reservation_release", + "budget_wait_scope": "github-default", + "budget_wait_revision": revision, + }) + monkeypatch.setattr( + worker, "load_providers", + lambda: (_ for _ in ()).throw( + AssertionError("deferred task loaded a provider"))) + + worker._execute_task_inner(waiter) + + deferred = isolated_db.get_task(created.id) + assert deferred.status.value == "pending" + assert deferred.scheduled_at == target + assert deferred.next_run_at is None + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT budget_wait_scope FROM tasks WHERE id = ?", + (created.id,), + ).fetchone() + assert stored["budget_wait_scope"] == "github-default" + + released_at = datetime.now(timezone.utc) + assert isolated_db.release_pipeline_github_budget( + "github-default", token="wake-owner", task_id=owner.id, + task_started_at=owner.started_at) is True + + woken = isolated_db.get_task(created.id) + assert woken.scheduled_at >= released_at + assert woken.scheduled_at <= datetime.now(timezone.utc) + assert woken.next_run_at is None + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT budget_wait_scope FROM tasks WHERE id = ?", + (created.id,), + ).fetchone() + assert stored["budget_wait_scope"] == "github-default" + + +def test_worker_preserves_budget_handoff_for_soft_retry( + isolated_db, monkeypatch): + created = isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + task = isolated_db.get_next_runnable() + revision = isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] + target = datetime.now(timezone.utc) + timedelta(seconds=5) + monkeypatch.setattr(pipeline_insights, "dispatch_gate", lambda _task: None) + monkeypatch.setattr( + pipeline_insights, "execution_route", lambda *_args, **_kwargs: { + "action": "defer", "mode": "skill", "reason": "ordered retry", + "defer_until": target.isoformat(), "defer_policy": "retry", + "budget_wait_scope": "github-default", + "budget_wait_revision": revision, + }) + monkeypatch.setattr( + worker, "load_providers", + lambda: (_ for _ in ()).throw( + AssertionError("deferred task loaded a provider"))) + + worker._execute_task_inner(task) + + deferred = isolated_db.get_task(created.id) + assert deferred.status.value == "pending" + assert target <= deferred.scheduled_at <= target + timedelta(seconds=1) + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT budget_wait_scope FROM tasks WHERE id = ?", + (created.id,), + ).fetchone() + assert stored["budget_wait_scope"] == "github-default" + + +def test_release_before_defer_cannot_lose_the_only_budget_wake( + isolated_db, monkeypatch): + owner = _running_task(isolated_db, "Budget owner") + assert _reserve( + isolated_db, owner, "racing-owner", core=500, + remaining=1000)["allowed"] is True + observed_revision = isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] + created = isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + assert waiter is not None and waiter.id == created.id + + # The owner finishes after admission observed contention but before the + # worker persists its defer. No waiter existed at release time. + assert isolated_db.release_pipeline_github_budget( + "github-default", token="racing-owner", task_id=owner.id, + task_started_at=owner.started_at) is True + target = datetime.now(timezone.utc) + timedelta(minutes=7) + monkeypatch.setattr(pipeline_insights, "dispatch_gate", lambda _task: None) + monkeypatch.setattr( + pipeline_insights, "execution_route", lambda *_args, **_kwargs: { + "action": "defer", "mode": "skill", "reason": "contention", + "defer_until": target.isoformat(), + "defer_policy": "reservation_release", + "budget_wait_scope": "github-default", + "budget_wait_revision": observed_revision, + }) + monkeypatch.setattr( + worker, "load_providers", + lambda: (_ for _ in ()).throw( + AssertionError("deferred task loaded a provider"))) + + before = datetime.now(timezone.utc) + worker._execute_task_inner(waiter) + + woken = isolated_db.get_task(created.id) + assert woken.status.value == "pending" + assert before <= woken.scheduled_at <= datetime.now(timezone.utc) + assert woken.next_run_at is None + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT budget_wait_scope FROM tasks WHERE id = ?", + (created.id,), + ).fetchone() + assert stored["budget_wait_scope"] == "github-default" + + +def test_budget_wait_revision_closes_release_and_reacquire_aba(isolated_db): + owner_a = _running_task(isolated_db, "Budget owner A") + owner_b = _running_task(isolated_db, "Budget owner B") + waiter = _running_task(isolated_db, "Example - REVIEW") + first = _reserve( + isolated_db, owner_a, "aba-owner-a", core=200, remaining=1000) + assert first["allowed"] is True + assert first["revision"] == 1 + + assert isolated_db.release_pipeline_github_budget( + "github-default", token="aba-owner-a", task_id=owner_a.id, + task_started_at=owner_a.started_at) is True + assert isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] == 2 + second = _reserve( + isolated_db, owner_b, "aba-owner-b", core=200, remaining=1000) + assert second["allowed"] is True + assert second["revision"] == 3 + + before = datetime.now(timezone.utc) + assert isolated_db.defer_task( + waiter.id, before + timedelta(minutes=7), "stale contention", + budget_wait_scope="github-default", + budget_wait_revision=first["revision"], + expected_started_at=waiter.started_at) is True + + retried = isolated_db.get_task(waiter.id) + assert before <= retried.scheduled_at <= datetime.now(timezone.utc) + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT budget_wait_scope FROM tasks WHERE id = ?", + (waiter.id,), + ).fetchone() + assert stored["budget_wait_scope"] == "github-default" + + +def test_released_budget_handoff_survives_claim_until_waiter_reserves( + isolated_db): + owner = _running_task(isolated_db, "Budget owner") + reserved = _reserve( + isolated_db, owner, "handoff-owner", core=500, remaining=1000) + assert reserved["allowed"] is True + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + assert isolated_db.defer_task( + waiter.id, datetime.now(timezone.utc) + timedelta(minutes=7), + "waiting for owner", budget_wait_scope="github-default", + budget_wait_revision=reserved["revision"], + expected_started_at=waiter.started_at) is True + low_created = isolated_db.create_task(TaskCreate( + prompt="Example - TAIL", recurrence="4h", priority=5)) + low = isolated_db.get_next_runnable() + assert low is not None and low.id == low_created.id + + assert isolated_db.release_pipeline_github_budget( + "github-default", token="handoff-owner", task_id=owner.id, + task_started_at=owner.started_at) is True + claimed_waiter = isolated_db.get_next_runnable() + assert claimed_waiter is not None and claimed_waiter.id == waiter.id + + denied = _reserve( + isolated_db, low, "lower-after-release", core=200, remaining=1000) + assert denied["allowed"] is False + assert denied["state"] == "priority_waiter" + + admitted = _reserve( + isolated_db, claimed_waiter, "waiter-after-release", + core=200, remaining=1000) + assert admitted["allowed"] is True + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT budget_wait_scope FROM tasks WHERE id = ?", + (waiter.id,), + ).fetchone() + assert stored["budget_wait_scope"] is None + + +def test_initial_budget_denial_arms_handoff_before_worker_defer( + isolated_db, monkeypatch): + execution = {"mode": "auto", "command": ["pipeline-tool"]} + profile = _profile_with_costs(execution=execution, core=200) + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "_github_rate_limits", lambda: _limits(core=1000)) + monkeypatch.setattr( + pipeline_insights, "_tool_available", + lambda *_args, **_kwargs: (_ for _ in ()).throw( + AssertionError("budget-denied task reached project preflight"))) + + owner = _running_task(isolated_db, "Budget owner") + assert _reserve( + isolated_db, owner, "arm-owner", core=800, + remaining=1000)["allowed"] is True + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + + route = pipeline_insights.execution_route( + waiter, waiter.prompt, retain_budget=True) + + assert route["action"] == "defer" + assert route["defer_policy"] == "reservation_release" + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT budget_wait_scope FROM tasks WHERE id = ?", + (waiter.id,), + ).fetchone() + assert stored["budget_wait_scope"] == "github-default" + + low = _running_task(isolated_db, "Example - TAIL") + assert isolated_db.release_pipeline_github_budget( + "github-default", token="arm-owner", task_id=owner.id, + task_started_at=owner.started_at) is True + denied = _reserve( + isolated_db, low, "lower-before-waiter-defer", + core=200, remaining=1000) + assert denied["allowed"] is False + assert denied["state"] == "priority_waiter" + + +def test_explicit_reschedule_releases_budget_priority_handoff(isolated_db): + waiter = _running_task(isolated_db, "Example - REVIEW") + revision = isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] + assert isolated_db.defer_task( + waiter.id, datetime.now(timezone.utc) + timedelta(minutes=7), + "waiting for owner", budget_wait_scope="github-default", + budget_wait_revision=revision, + expected_started_at=waiter.started_at) is True + + assert isolated_db.update_task_fields( + waiter.id, + {"scheduled_at": datetime.now(timezone.utc) + timedelta(hours=3)}) is True + + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT budget_wait_scope FROM tasks WHERE id = ?", + (waiter.id,), + ).fetchone() + assert stored["budget_wait_scope"] is None + + +def test_stale_prune_yields_to_woken_higher_priority_waiter(isolated_db): + owner = _running_task(isolated_db, "Crashed budget owner") + reserved = _reserve( + isolated_db, owner, "stale-owner", core=500, remaining=1000) + assert reserved["allowed"] is True + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + assert isolated_db.defer_task( + waiter.id, datetime.now(timezone.utc) + timedelta(minutes=7), + "waiting for stale owner", budget_wait_scope="github-default", + budget_wait_revision=reserved["revision"], + expected_started_at=waiter.started_at) is True + low_created = isolated_db.create_task(TaskCreate( + prompt="Example - TAIL", recurrence="4h", priority=5)) + low = isolated_db.get_next_runnable() + assert low is not None and low.id == low_created.id + + # Simulate terminal state committed before explicit reservation cleanup. + assert isolated_db.mark_completed( + owner.id, "done", expected_started_at=owner.started_at) is True + denied = _reserve( + isolated_db, low, "lower-priority", core=200, remaining=1000) + + assert denied["allowed"] is False + assert denied["state"] == "priority_waiter" + assert isolated_db.pipeline_github_budget_reservations( + "github-default")["count"] == 0 + woken = isolated_db.get_task(waiter.id) + assert woken.scheduled_at <= datetime.now(timezone.utc) + assert woken.next_run_at is None + + +def test_execution_route_yields_when_initial_scan_wakes_priority_waiter( + isolated_db, monkeypatch): + execution = {"mode": "auto", "command": ["pipeline-tool"]} + profile = _profile_with_costs(execution=execution, core=200) + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "_github_rate_limits", lambda: _limits(core=1000)) + monkeypatch.setattr( + pipeline_insights, "_tool_available", + lambda *_args, **_kwargs: (_ for _ in ()).throw( + AssertionError("lower-priority task reached project preflight"))) + + owner = _running_task(isolated_db, "Crashed budget owner") + reserved = _reserve( + isolated_db, owner, "scan-stale-owner", core=500, remaining=1000) + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + assert isolated_db.defer_task( + waiter.id, datetime.now(timezone.utc) + timedelta(minutes=7), + "waiting for stale owner", budget_wait_scope="github-default", + budget_wait_revision=reserved["revision"], + expected_started_at=waiter.started_at) is True + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=5)) + lower = isolated_db.get_next_runnable() + assert isolated_db.mark_completed( + owner.id, "done", expected_started_at=owner.started_at) is True + + route = pipeline_insights.execution_route( + lower, lower.prompt, retain_budget=True) + + assert route["action"] == "defer" + assert route["defer_policy"] == "retry" + assert route["github_budget"]["state"] == "priority_waiter" + assert isolated_db.pipeline_github_budget_reservations( + "github-default")["count"] == 0 + woken = isolated_db.get_task(waiter.id) + assert woken.scheduled_at <= datetime.now(timezone.utc) + assert woken.next_run_at is None + + +def test_execution_route_skips_preflight_for_existing_priority_waiter( + isolated_db, monkeypatch): + execution = {"mode": "auto", "command": ["pipeline-tool"]} + profile = _profile_with_costs(execution=execution, core=200) + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "_github_rate_limits", lambda: _limits(core=1000)) + monkeypatch.setattr( + pipeline_insights, "_tool_available", + lambda *_args, **_kwargs: (_ for _ in ()).throw( + AssertionError("lower-priority task reached project preflight"))) + + owner = _running_task(isolated_db, "Budget owner") + reserved = _reserve( + isolated_db, owner, "live-owner", core=200, remaining=1000) + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + assert isolated_db.defer_task( + waiter.id, datetime.now(timezone.utc) + timedelta(minutes=7), + "waiting for live owner", budget_wait_scope="github-default", + budget_wait_revision=reserved["revision"], + expected_started_at=waiter.started_at) is True + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=5)) + lower = isolated_db.get_next_runnable() + + route = pipeline_insights.execution_route( + lower, lower.prompt, retain_budget=True) + + assert route["action"] == "defer" + assert route["defer_policy"] == "retry" + assert route["github_budget"]["state"] == "priority_waiter" + ledger = isolated_db.pipeline_github_budget_reservations( + "github-default") + assert ledger["count"] == 1 + assert ledger["items"][0]["token"] == "live-owner" + + +def test_final_reserve_cas_retries_when_release_happens_during_preflight( + isolated_db, monkeypatch): + execution = {"mode": "auto", "command": ["pipeline-tool"]} + profile = _profile_with_costs(execution=execution, core=200) + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "_github_rate_limits", lambda: _limits(core=1000)) + monkeypatch.setattr( + pipeline_insights, "_tool_available", lambda *_args: (True, "")) + + owner = _running_task(isolated_db, "Budget owner") + reserved = _reserve( + isolated_db, owner, "preflight-owner", core=200, remaining=1000) + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=5)) + lower = isolated_db.get_next_runnable() + + def release_during_preflight(*_args, **_kwargs): + assert isolated_db.defer_task( + waiter.id, datetime.now(timezone.utc) + timedelta(minutes=7), + "waiting for owner", budget_wait_scope="github-default", + budget_wait_revision=reserved["revision"], + expected_started_at=waiter.started_at) is True + assert isolated_db.release_pipeline_github_budget( + "github-default", token="preflight-owner", task_id=owner.id, + task_started_at=owner.started_at) is True + return {"action": "fallback", "reason": "full review required"} + + monkeypatch.setattr( + pipeline_insights, "_tool_preflight", release_during_preflight) + + route = pipeline_insights.execution_route( + lower, lower.prompt, retain_budget=True) + + assert route["action"] == "defer" + assert route["defer_policy"] == "retry" + assert route["github_budget"]["state"] == "ledger_changed" + assert isolated_db.pipeline_github_budget_reservations( + "github-default")["count"] == 0 + woken = isolated_db.get_task(waiter.id) + assert woken.scheduled_at <= datetime.now(timezone.utc) + assert woken.next_run_at is None + + +def test_ledger_changed_retry_preserves_handoff_until_reservation( + isolated_db, monkeypatch): + execution = {"mode": "auto", "command": ["pipeline-tool"]} + profile = _profile_with_costs(execution=execution, core=200) + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "_github_rate_limits", lambda: _limits(core=1000)) + monkeypatch.setattr( + pipeline_insights, "_tool_available", lambda *_args: (True, "")) + + owner = _running_task(isolated_db, "Budget owner") + reserved = _reserve( + isolated_db, owner, "cas-handoff-owner", core=200, remaining=1000) + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + armed = isolated_db.arm_pipeline_github_budget_waiter( + "github-default", task_id=waiter.id, + task_started_at=waiter.started_at, + expected_revision=reserved["revision"]) + assert armed["armed"] is True + low_created = isolated_db.create_task(TaskCreate( + prompt="Example - TAIL", recurrence="4h", priority=5)) + low = isolated_db.get_next_runnable() + assert low is not None and low.id == low_created.id + + def release_during_preflight(*_args, **_kwargs): + assert isolated_db.release_pipeline_github_budget( + "github-default", token="cas-handoff-owner", task_id=owner.id, + task_started_at=owner.started_at) is True + return {"action": "fallback", "reason": "full review required"} + + monkeypatch.setattr( + pipeline_insights, "_tool_preflight", release_during_preflight) + + route = pipeline_insights.execution_route( + waiter, waiter.prompt, retain_budget=True) + + assert route["action"] == "defer" + assert route["defer_policy"] == "retry" + assert route["github_budget"]["state"] == "ledger_changed" + assert route["budget_wait_scope"] == "github-default" + assert type(route["budget_wait_revision"]) is int + assert isolated_db.defer_task( + waiter.id, datetime.fromisoformat(route["defer_until"]), + route["reason"], budget_wait_scope=route["budget_wait_scope"], + budget_wait_revision=route["budget_wait_revision"], + expected_started_at=waiter.started_at) is True + + denied = _reserve( + isolated_db, low, "lower-after-cas-retry", + core=200, remaining=1000) + assert denied["allowed"] is False + assert denied["state"] == "priority_waiter" + + +def test_priority_waiter_retry_preserves_order_behind_higher_waiter( + isolated_db, monkeypatch): + execution = {"mode": "auto", "command": ["pipeline-tool"]} + profile = _profile_with_costs(execution=execution, core=200) + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "_github_rate_limits", lambda: _limits(core=1000)) + monkeypatch.setattr( + pipeline_insights, "_tool_available", + lambda *_args, **_kwargs: (_ for _ in ()).throw( + AssertionError("priority waiter reached project preflight"))) + + isolated_db.create_task(TaskCreate( + prompt="Example - MERGE", recurrence="4h", priority=1)) + first = isolated_db.get_next_runnable() + revision = isolated_db.pipeline_github_budget_reservations( + "github-default")["revision"] + assert isolated_db.defer_task( + first.id, datetime.now(timezone.utc) + timedelta(minutes=7), + "first waiter", budget_wait_scope="github-default", + budget_wait_revision=revision, + expected_started_at=first.started_at) is True + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + second = isolated_db.get_next_runnable() + + route = pipeline_insights.execution_route( + second, second.prompt, retain_budget=True) + + assert route["action"] == "defer" + assert route["defer_policy"] == "retry" + assert route["github_budget"]["state"] == "priority_waiter" + assert route["budget_wait_scope"] == "github-default" + assert isolated_db.defer_task( + second.id, datetime.fromisoformat(route["defer_until"]), + route["reason"], budget_wait_scope=route["budget_wait_scope"], + budget_wait_revision=route["budget_wait_revision"], + expected_started_at=second.started_at) is True + assert isolated_db.cancel_task(first.id) is True + + isolated_db.create_task(TaskCreate( + prompt="Example - FIX", recurrence="4h", priority=3)) + third = isolated_db.get_next_runnable() + denied = _reserve( + isolated_db, third, "third-behind-second", core=200, remaining=1000) + assert denied["allowed"] is False + assert denied["state"] == "priority_waiter" + + +def test_scan_lease_busy_retry_preserves_existing_priority_handoff( + isolated_db, monkeypatch): + execution = {"mode": "auto", "command": ["pipeline-tool"]} + profile = _profile_with_costs(execution=execution, core=200) + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "_github_rate_limits", lambda: _limits(core=1000)) + monkeypatch.setattr( + pipeline_insights, "_tool_available", + lambda *_args, **_kwargs: (_ for _ in ()).throw( + AssertionError("lower-priority task reached project preflight"))) + + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + armed = isolated_db.arm_pipeline_github_budget_waiter( + "github-default", task_id=waiter.id, + task_started_at=waiter.started_at) + assert armed["armed"] is True + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=5)) + lower = isolated_db.get_next_runnable() + + held = isolated_db.acquire_pipeline_scan_lease( + "github-default", "lower-held-scan", 30) + assert held["acquired"] is True + + route = pipeline_insights.execution_route( + waiter, waiter.prompt, retain_budget=True) + + assert route["action"] == "defer" + assert route["defer_policy"] == "retry" + assert route["github_budget"]["state"] == "scan_in_progress" + assert route["budget_wait_scope"] == "github-default" + assert isolated_db.defer_task( + waiter.id, datetime.fromisoformat(route["defer_until"]), + route["reason"], budget_wait_scope=route["budget_wait_scope"], + budget_wait_revision=route["budget_wait_revision"], + expected_started_at=waiter.started_at) is True + assert isolated_db.release_pipeline_scan_lease( + "github-default", "lower-held-scan") is True + + lower_route = pipeline_insights.execution_route( + lower, lower.prompt, retain_budget=True) + assert lower_route["action"] == "defer" + assert lower_route["github_budget"]["state"] == "priority_waiter" + + +def test_safe_corrupt_ledger_recovery_advances_revision_and_wakes_waiter( + isolated_db): + owner = _running_task(isolated_db, "Budget owner") + reserved = _reserve( + isolated_db, owner, "corrupt-owner", core=200, remaining=1000) + waiter = _running_task(isolated_db, "Example - REVIEW") + target = datetime.now(timezone.utc) + timedelta(minutes=7) + assert isolated_db.defer_task( + waiter.id, target, "waiting for owner", + budget_wait_scope="github-default", + budget_wait_revision=reserved["revision"], + expected_started_at=waiter.started_at) is True + assert isolated_db.mark_completed( + owner.id, "done", expected_started_at=owner.started_at) is True + with isolated_db._connect() as conn: + conn.execute( + "INSERT OR REPLACE INTO settings (key, value) VALUES (?, ?)", + (isolated_db._pipeline_github_budget_reservation_key( + "github-default"), "not-json"), + ) + + ledger = isolated_db.pipeline_github_budget_reservations( + "github-default") + + assert ledger["count"] == 0 + assert ledger["revision"] == reserved["revision"] + 1 + woken = isolated_db.get_task(waiter.id) + assert woken.scheduled_at <= datetime.now(timezone.utc) + assert woken.next_run_at is None + + +def test_live_rate_reset_deadline_is_not_woken_by_reservation_release( + isolated_db): + owner = _running_task(isolated_db, "Budget owner") + assert _reserve( + isolated_db, owner, "hard-deadline-owner", core=200, + remaining=1000)["allowed"] is True + created = isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=2)) + waiter = isolated_db.get_next_runnable() + target = datetime.now(timezone.utc) + timedelta(minutes=7) + assert isolated_db.defer_task( + waiter.id, target, "live rate limit", hard_not_before=True, + expected_started_at=waiter.started_at) is True + + assert isolated_db.release_pipeline_github_budget( + "github-default", token="hard-deadline-owner", task_id=owner.id, + task_started_at=owner.started_at) is True + + deferred = isolated_db.get_task(created.id) + assert deferred.scheduled_at == target + assert deferred.next_run_at == target + + def test_worker_low_budget_stays_deferred_when_status_database_is_locked( isolated_db, monkeypatch): created = isolated_db.create_task(TaskCreate(