From aa1b32db75e770d274b76900365a98e8ac963625 Mon Sep 17 00:00:00 2001 From: ibrog Date: Thu, 17 Sep 2026 18:58:51 +0300 Subject: [PATCH] fix(pipeline): unblock integration review budget Let a signed, task-fenced integration review occurrence borrow the active P1 merge budget priority without promoting ordinary review work. Persist and safely restore the occurrence priority across budget deferrals, while preserving explicit operator edits and recurring-series priority. Generated-with: Codex --- README.md | 9 ++ promptpilot/db.py | 89 ++++++++++- promptpilot/pipeline_insights.py | 159 +++++++++++++++++--- tests/test_github_budget.py | 71 +++++++++ tests/test_pipeline_replicas.py | 247 +++++++++++++++++++++++++++++++ 5 files changed, 555 insertions(+), 20 deletions(-) diff --git a/README.md b/README.md index 15ca910..d6fda8d 100644 --- a/README.md +++ b/README.md @@ -1217,6 +1217,15 @@ REST health не доказывает timeline lineage: одинаковый с устаревший `ship`. Поэтому такой снимок нельзя превращать в `action=wait` без полного fallback-гейта. +Если для MERGE настроен более высокий приоритет и `priority_one_headroom`, +подписанный и привязанный к точной попытке `integration-review` временно +наследует приоритет активной MERGE-серии: этот REVIEW является зависимостью, +которая освобождает сам MERGE. Обычный content REVIEW запас не использует. +Временный приоритет хранится только у текущего вхождения серии, переживает +смену попытки при безопасном +budget-defer, снимается при смене маршрута, а следующий запуск повторяющейся +REVIEW-серии снова получает её настроенный приоритет. + Встроенный `promptpilot.project_pipeline` реализует общий protocol v1. Проект задаёт `repository`, доверенный аккаунт, health-команду, base branch и required checks в `pipelinectl.json`. Обычный REVIEW получает ровно один кандидат и diff --git a/promptpilot/db.py b/promptpilot/db.py index 9273793..69f7355 100644 --- a/promptpilot/db.py +++ b/promptpilot/db.py @@ -72,6 +72,7 @@ worktree_branch TEXT, note TEXT, budget_wait_scope TEXT, + pipeline_priority_restore INTEGER, verdict TEXT ,series_id INTEGER REFERENCES task_series(id) ); @@ -393,6 +394,7 @@ "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)", + "ALTER TABLE tasks ADD COLUMN pipeline_priority_restore INTEGER", ] WORKFLOW_SCHEMA_VERSION = "workflow_orchestrator_w0_v1" @@ -429,9 +431,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. + # 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("pipeline_priority_restore", None) for field in ("scheduled_at", "next_run_at", "created_at", "started_at", "completed_at"): d[field] = _parse_dt(d[field]) return TaskInDB(**d) @@ -1479,12 +1482,84 @@ 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 = ? WHERE id = ? AND status IN ('pending', 'rate_limited')", + """UPDATE tasks SET priority = ?, pipeline_priority_restore = NULL + WHERE id = ? AND status IN ('pending', 'rate_limited')""", (priority, task_id), ) return cur.rowcount > 0 +def promote_running_attempt_priority(task_id: int, started_at, + priority: int) -> int | None: + """Promote one exact running attempt without changing its series. + + A routed pipeline occurrence can become more urgent only after its signed + target has been inspected. Keep that temporary promotion on the task row + so budget-wait ordering observes it, while the recurring successor still + inherits the configured series priority. + """ + if (isinstance(task_id, bool) or not isinstance(task_id, int) or task_id <= 0 + or isinstance(priority, bool) or not isinstance(priority, int) + or not 1 <= priority <= 10): + raise ValueError("invalid running task priority promotion") + attempt = _attempt_iso(started_at) + if attempt is None: + raise ValueError("running task attempt is required for priority promotion") + with _connect(immediate=True) as conn: + row = conn.execute( + """SELECT priority, pipeline_priority_restore FROM tasks + WHERE id = ? AND status = 'running' AND started_at = ?""", + (task_id, attempt), + ).fetchone() + if row is None: + return None + current = int(row["priority"]) + restore = row["pipeline_priority_restore"] + if restore is not None and not 1 <= int(restore) <= 10: + raise ValueError("stored pipeline priority restore value is invalid") + if current > priority: + conn.execute( + """UPDATE tasks + SET priority = ?, pipeline_priority_restore = COALESCE( + pipeline_priority_restore, priority) + WHERE id = ? AND status = 'running' AND started_at = ? + AND priority > ?""", + (priority, task_id, attempt, priority), + ) + return priority + return current + + +def restore_running_attempt_priority(task_id: int, started_at) -> int | None: + """Drop a temporary routed promotion after fresh preflight disproves it.""" + if isinstance(task_id, bool) or not isinstance(task_id, int) or task_id <= 0: + raise ValueError("invalid running task priority restoration") + attempt = _attempt_iso(started_at) + if attempt is None: + raise ValueError("running task attempt is required for priority restoration") + with _connect(immediate=True) as conn: + row = conn.execute( + """SELECT priority, pipeline_priority_restore FROM tasks + WHERE id = ? AND status = 'running' AND started_at = ?""", + (task_id, attempt), + ).fetchone() + if row is None: + return None + restore = row["pipeline_priority_restore"] + if restore is None: + return int(row["priority"]) + restored = int(restore) + if not 1 <= restored <= 10: + raise ValueError("stored pipeline priority restore value is invalid") + conn.execute( + """UPDATE tasks + SET priority = ?, pipeline_priority_restore = NULL + WHERE id = ? AND status = 'running' AND started_at = ?""", + (restored, task_id, attempt), + ) + return restored + + def _effective_series_recurrence(row: dict, now: Optional[datetime] = None) -> str: now = now or datetime.now(timezone.utc) temporary = row.get("temporary_recurrence") @@ -1653,6 +1728,10 @@ def update_series(series_id: int, fields: dict) -> bool: task_fields[name] = fields[name] task_sets = [f"{k} = ?" for k in task_fields] 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") 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 @@ -2347,6 +2426,10 @@ def update_task_fields(task_id: int, fields: dict) -> bool: # 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 + if "priority" in fields: + # Do not let a later route restoration undo the operator's explicit + # priority choice for a temporarily promoted waiting occurrence. + fields["pipeline_priority_restore"] = None sets = ", ".join(f"{k} = ?" for k in fields) with _connect() as conn: cur = conn.execute( diff --git a/promptpilot/pipeline_insights.py b/promptpilot/pipeline_insights.py index 56993d1..4a1fb5d 100644 --- a/promptpilot/pipeline_insights.py +++ b/promptpilot/pipeline_insights.py @@ -423,6 +423,46 @@ def _budget_policy_for_admission_priority(policy: dict, return selected +def _post_preflight_admission_priority( + profile: dict, series: list[dict], stage: str, + provider_route: str | None, + validated_target_stage: str | None, + effective_priority_headroom: dict | None) -> int | None: + """Let the exact integration REVIEW use the MERGE headroom it unlocks. + + ``validated_target_stage`` is populated only from the signed, task-fenced + replica lease. Ordinary content review remains at the configured series + priority and therefore continues to preserve priority-1 capacity. + """ + from .fallback_handoff import INTEGRATION_REVIEW_STAGES + + if (stage != "review" or provider_route != "fallback_targeted" + or validated_target_stage not in INTEGRATION_REVIEW_STAGES): + return None + # Admission has already combined every profile that shares this GitHub + # account. Looking only at the current profile would reintroduce the cycle + # whenever another repository owns the strongest shared headroom. + headroom = effective_priority_headroom or {} + if (not isinstance(headroom, dict) + or not any(type(value) is int and value > 0 + for value in headroom.values())): + return None + priorities = [] + for queue in profile.get("queues", []): + execution = queue.get("execution") or {} + queue_stage = str(execution.get("stage") or queue.get("id") or "").lower() + if queue_stage != "merge": + continue + for item in _series_replicas_for_queue(queue, series): + priority = item.get("priority") + if (not item.get("paused") and not item.get("broken") + and type(priority) is int and 1 <= priority <= 10): + priorities.append(priority) + # This reserve is specifically consumable by priority 1. A lower-urgency + # MERGE series does not justify leaving the REVIEW occurrence promoted. + return 1 if 1 in priorities else None + + def _budget_policy_for_route(policy: dict, route: str | None) -> dict: """Attach an estimated route cost; profiles without costs stay legacy.""" if route is None or policy.get("costs") is None: @@ -1005,10 +1045,14 @@ def _reservation_denied(policy: dict, limits: dict | None, result: dict, *, def _reserve_execution_admission( task, profile_id: str, profile: dict, queue: dict, admission: dict, - budget_route: str, *, retain_budget: bool) -> dict: + budget_route: str, *, retain_budget: bool, + admission_priority: int | None = None) -> dict: """Recheck the elected route and optionally reserve its in-flight cost.""" lease = admission.get("_lease") expected_revision = admission.get("_budget_reservation_revision") + effective_priority = ( + getattr(task, "priority", None) + if admission_priority is None else admission_priority) if type(expected_revision) is not int or expected_revision < 0: expected_revision = None try: @@ -1018,7 +1062,7 @@ def _reserve_execution_admission( policy = _with_shared_budget_floor(policy) policy = _budget_policy_for_route(policy, budget_route) policy = _budget_policy_for_admission_priority( - policy, getattr(task, "priority", None)) + policy, effective_priority) 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. @@ -1047,7 +1091,7 @@ def _reserve_execution_admission( policy["lease_scope"]) decision = _priority_waiter_decision( policy, limits, reservations, - priority=getattr(task, "priority", None), + priority=effective_priority, status_revision=lease.status_revision) if decision is None: decision = _evaluate_github_budget( @@ -2689,8 +2733,9 @@ def execution_route(task, fallback_prompt: str, working_dir: str | None = None, "profile_id": profile_id, "queue_id": queue.get("id"), } replicated = replica_count > 1 + profile_series = db.list_series() if replicated else [] if replicated: - replica_status = _queue_replica_status(queue, db.list_series()) + replica_status = _queue_replica_status(queue, profile_series) stage_name = str((queue.get("execution") or {}).get("stage") or queue.get("id") or "").lower() issues = list(replica_status["issues"]) @@ -2889,6 +2934,8 @@ def execution_route(task, fallback_prompt: str, working_dir: str | None = None, preflight_action = preflight["action"].lower() target_reservation = None + validated_target_stage = None + fallback_lease = None if replicated and preflight_action in {"audit", "fallback"}: try: from . import project_pipeline @@ -2925,6 +2972,7 @@ def execution_route(task, fallback_prompt: str, working_dir: str | None = None, raise project_pipeline.PipelineError( "pipeline target reservation contradicts task, repository, or target") target_reservation = reservation + validated_target_stage = target_stage except (project_pipeline.PipelineError, TypeError, ValueError) as exc: return { "action": "block", "mode": "tool", @@ -2932,6 +2980,37 @@ def execution_route(task, fallback_prompt: str, working_dir: str | None = None, "profile_id": profile_id, "queue_id": queue.get("id"), "preflight": preflight, } + if (preflight_action == "fallback" and "handoff" in preflight + and target_reservation is not None): + from .fallback_handoff import validate + from .project_pipeline import PipelineError + + try: + fallback_lease = validate(preflight, stage) + if command[-2:] != ["next", stage]: + raise PipelineError( + "fallback handoff command must end with next and the exact stage") + except (PipelineError, TypeError, ValueError) as exc: + try: + task_id = getattr(task, "id", None) + started_at = getattr(task, "started_at", None) + if type(task_id) is int and task_id > 0 and started_at is not None: + db.restore_running_attempt_priority(task_id, started_at) + except (sqlite3.Error, TypeError, ValueError) as restore_exc: + denied = _replace_admission_with_lease_failure( + admission, profile, _GitHubScanLeaseUnavailable( + "pipeline admission priority restoration is unavailable " + f"after invalid fallback handoff: {restore_exc}")) + return _budget_defer_route( + denied, profile_id, profile, queue, mode="skill", + phase="post_preflight", preflight=preflight) + return { + "action": "block", "mode": "tool", "reason": str(exc), + "profile_id": profile_id, "queue_id": queue.get("id"), + "preflight": preflight, + } + validated_target_stage = str( + (fallback_lease.get("target") or {}).get("stage") or "").lower() provider_route = None if preflight_action == "fallback" and mode == "auto": provider_route = ("fallback_targeted" if "handoff" in preflight @@ -2940,10 +3019,54 @@ def execution_route(task, fallback_prompt: str, working_dir: str | None = None, provider_route = "tool" elif preflight_action not in {"empty", "wait", "error"} and mode == "auto": provider_route = "skill" + routed_priority = _post_preflight_admission_priority( + profile, profile_series, stage, provider_route, + validated_target_stage, admission.get("priority_one_headroom")) + # Only a signed target that is also fenced to this exact replica task + # may borrow the priority-1 integration headroom. + if routed_priority is not None and target_reservation is None: + routed_priority = None + task_id = getattr(task, "id", None) + task_started_at = getattr(task, "started_at", None) + admission_priority = getattr(task, "priority", None) + if type(task_id) is int and task_id > 0 and task_started_at is not None: + try: + if routed_priority is None: + admission_priority = db.restore_running_attempt_priority( + task_id, task_started_at) + else: + admission_priority = db.promote_running_attempt_priority( + task_id, task_started_at, routed_priority) + except (sqlite3.Error, TypeError, ValueError) as exc: + denied = _replace_admission_with_lease_failure( + admission, profile, _GitHubScanLeaseUnavailable( + f"pipeline admission priority transition is unavailable: {exc}")) + return _budget_defer_route( + denied, profile_id, profile, queue, + mode=("tool" if provider_route == "tool" else "skill"), + phase="post_preflight", preflight=preflight) + if admission_priority is None: + return { + "action": "block", "mode": "tool", + "reason": ( + "running task attempt changed before pipeline admission " + "priority transition"), + "profile_id": profile_id, "queue_id": queue.get("id"), + "preflight": preflight, + } + elif routed_priority is not None: + return { + "action": "block", "mode": "tool", + "reason": ( + "integration REVIEW admission has no exact running task attempt"), + "profile_id": profile_id, "queue_id": queue.get("id"), + "preflight": preflight, + } if provider_route is not None: admission = _reserve_execution_admission( task, profile_id, profile, queue, admission, provider_route, - retain_budget=retain_budget) + retain_budget=retain_budget, + admission_priority=admission_priority) if not admission.get("allowed"): return _budget_defer_route( admission, profile_id, profile, queue, @@ -2977,19 +3100,21 @@ def execution_route(task, fallback_prompt: str, working_dir: str | None = None, } if preflight_action == "fallback": if "handoff" in preflight: - from .fallback_handoff import validate - from .project_pipeline import PipelineError + if fallback_lease is None: + from .fallback_handoff import validate + from .project_pipeline import PipelineError - try: - fallback_lease = validate(preflight, stage) - if command[-2:] != ["next", stage]: - raise PipelineError("fallback handoff command must end with next and the exact stage") - except (PipelineError, TypeError, ValueError) as exc: - return { - "action": "block", "mode": "tool", "reason": str(exc), - "profile_id": profile_id, "queue_id": queue.get("id"), - "preflight": preflight, - } + try: + fallback_lease = validate(preflight, stage) + if command[-2:] != ["next", stage]: + raise PipelineError( + "fallback handoff command must end with next and the exact stage") + except (PipelineError, TypeError, ValueError) as exc: + return { + "action": "block", "mode": "tool", "reason": str(exc), + "profile_id": profile_id, "queue_id": queue.get("id"), + "preflight": preflight, + } if mode == "auto": gate_command = [*command[:-2], "gate-fallback", stage, "--lease", preflight["handoff"]["lease"]] diff --git a/tests/test_github_budget.py b/tests/test_github_budget.py index 86329a7..bae940c 100644 --- a/tests/test_github_budget.py +++ b/tests/test_github_budget.py @@ -198,6 +198,77 @@ def test_priority_one_headroom_defaults_to_zero_and_accepts_exact_vector(): } +def test_route_priority_promotion_is_attempt_fenced_internal_and_reversible( + isolated_db): + created = isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="15m", priority=2)) + task = isolated_db.get_next_runnable() + assert task.id == created.id + wrong_attempt = task.started_at + timedelta(microseconds=1) + + assert isolated_db.promote_running_attempt_priority( + task.id, wrong_attempt, 1) is None + assert isolated_db.get_task(task.id).priority == 2 + assert isolated_db.promote_running_attempt_priority( + task.id, task.started_at, 1) == 1 + assert isolated_db.promote_running_attempt_priority( + task.id, task.started_at, 1) == 1 + promoted = isolated_db.get_task(task.id) + assert promoted.priority == 1 + assert "pipeline_priority_restore" not in promoted.model_dump() + with isolated_db._connect() as conn: + row = conn.execute( + "SELECT pipeline_priority_restore FROM tasks WHERE id = ?", + (task.id,), + ).fetchone() + columns = {item["name"] for item in conn.execute( + "PRAGMA table_info(tasks)").fetchall()} + assert row["pipeline_priority_restore"] == 2 + assert "pipeline_priority_restore" in columns + + assert isolated_db.restore_running_attempt_priority( + task.id, wrong_attempt) is None + assert isolated_db.restore_running_attempt_priority( + task.id, task.started_at) == 2 + restored = isolated_db.get_task(task.id) + assert restored.priority == 2 + with isolated_db._connect() as conn: + row = conn.execute( + "SELECT pipeline_priority_restore FROM tasks WHERE id = ?", + (task.id,), + ).fetchone() + assert row["pipeline_priority_restore"] is None + + +@pytest.mark.parametrize("edit_scope", ["task", "series"]) +def test_explicit_priority_edit_supersedes_temporary_route_promotion( + isolated_db, edit_scope): + created = isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="15m", priority=2)) + task = isolated_db.get_next_runnable() + assert task.id == created.id + assert isolated_db.promote_running_attempt_priority( + task.id, task.started_at, 1) == 1 + assert isolated_db.defer_task( + task.id, datetime.now(timezone.utc) + timedelta(minutes=5), + "budget wait", expected_started_at=task.started_at) + + if edit_scope == "task": + assert isolated_db.update_task_fields(task.id, {"priority": 3}) + expected_priority = 3 + else: + assert isolated_db.update_series(task.series_id, {"priority": 4}) + expected_priority = 4 + + assert isolated_db.get_task(task.id).priority == expected_priority + with isolated_db._connect() as conn: + row = conn.execute( + "SELECT pipeline_priority_restore FROM tasks WHERE id = ?", + (task.id,), + ).fetchone() + assert row["pipeline_priority_restore"] is None + + @pytest.mark.parametrize("invalid", [ [], {"core": 1200, "search": 12}, diff --git a/tests/test_pipeline_replicas.py b/tests/test_pipeline_replicas.py index 1bb1abd..2d3b207 100644 --- a/tests/test_pipeline_replicas.py +++ b/tests/test_pipeline_replicas.py @@ -1004,6 +1004,253 @@ def preflight(_execution, _command, _working_dir, *, env_extra=None): } +def test_signed_integration_review_inherits_merge_budget_priority( + isolated_db, monkeypatch, tmp_path): + first_dir = tmp_path / "review-1" + second_dir = tmp_path / "review-2" + merge_dir = tmp_path / "merge" + for directory in (first_dir, second_dir, merge_dir): + directory.mkdir() + first = isolated_db.create_task(TaskCreate( + prompt="Project - REVIEW 1", working_dir=str(first_dir), + recurrence="15m", priority=2)) + isolated_db.create_task(TaskCreate( + prompt="Project - REVIEW 2", working_dir=str(second_dir), + recurrence="15m", priority=2)) + task = isolated_db.get_next_runnable() + assert task.id == first.id + merge = isolated_db.create_task(TaskCreate( + prompt="Project - MERGE", working_dir=str(merge_dir), recurrence="10m", + priority=1, scheduled_at=datetime.now(timezone.utc) + timedelta(hours=1))) + reservation = isolated_db.reserve_pipeline_target( + "owner/repo", "integration-review", 10, HEAD_A, task.id, 300, + task_started_at=task.started_at.astimezone(timezone.utc).isoformat(), + ownership_kind="headless") + review_queue = { + "id": "review", "series_contains": "Project - REVIEW", "replicas": 2, + "execution": { + "mode": "auto", "stage": "review", + "command": ["pipelinectl", "next", "review"], + }, + } + merge_queue = { + "id": "merge", "series_contains": "Project - MERGE", + "execution": {"mode": "auto", "stage": "merge"}, + } + profile = { + "repository": "owner/repo", "queues": [review_queue, merge_queue], + "github_budget": { + "minimum_remaining": {"core": 250, "search": 0, "graphql": 0}, + "priority_one_headroom": { + "core": 1200, "search": 0, "graphql": 0, + }, + "costs": { + "insights": {"core": 0, "search": 0, "graphql": 0}, + "skill": {"core": 1000, "search": 0, "graphql": 0}, + "tool_preflight": {"core": 600, "search": 0, "graphql": 0}, + "tool": {"core": 900, "search": 0, "graphql": 0}, + "fallback_targeted": { + "core": 1000, "search": 0, "graphql": 0, + }, + }, + }, + } + target = {"number": 10, "head": HEAD_A, "stage": "integration-review"} + lease = { + "repository": "owner/repo", "target": target, + "target_reservation": reservation, "pipeline_replicas": 2, + } + preflight = { + "action": "fallback", "reason": "integration owner", "target": target, + "handoff": {"lease": "signed-integration-lease"}, + } + reset = int(time.time()) + 3600 + # First election passes as P2 at 2399, but a live drop makes the final P1 + # reservation wait. After defer/reclaim, 1250 is enough only because the + # proven integration occurrence retained P1 for its new attempt. + core_remaining = iter((2399, 1100, 1250, 1250)) + + def limits(): + core = next(core_remaining) + return { + "core": {"limit": 5000, "used": 5000 - core, "remaining": core, + "reset": reset, "reset_at": None}, + "search": {"limit": 30, "used": 0, "remaining": 30, + "reset": reset, "reset_at": None}, + "graphql": {"limit": 5000, "used": 0, "remaining": 5000, + "reset": reset, "reset_at": None}, + } + monkeypatch.setattr( + pipeline_insights, "_matching_queue", + lambda _task: ("p", profile, review_queue)) + monkeypatch.setattr(pipeline_insights, "_profiles", lambda: {"p": profile}) + monkeypatch.setattr( + pipeline_insights, "_github_rate_limits", limits) + monkeypatch.setattr( + pipeline_insights, "_tool_available", lambda *_args: (True, "")) + monkeypatch.setattr( + pipeline_insights, "_tool_preflight", lambda *_args, **_kwargs: preflight) + monkeypatch.setattr(pipeline_insights, "load_providers", lambda: {}) + monkeypatch.setattr(pipelinectl, "decode_signed_lease", lambda _value: lease) + monkeypatch.setattr(fallback_handoff, "validate", lambda *_args: lease) + pipeline_insights.release_execution_admission() + + first_route = pipeline_insights.execution_route( + task, "skill", str(first_dir), retain_budget=True) + assert first_route["action"] == "defer" + assert isolated_db.get_task(task.id).priority == 1 + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT pipeline_priority_restore FROM tasks WHERE id = ?", + (task.id,), + ).fetchone() + assert stored["pipeline_priority_restore"] == 2 + pipeline_insights.release_execution_admission() + + assert isolated_db.defer_task( + task.id, datetime.now(timezone.utc) - timedelta(seconds=1), + first_route["reason"], expected_started_at=task.started_at) + reclaimed = isolated_db.get_next_runnable() + assert reclaimed.id == task.id + assert reclaimed.started_at != task.started_at + assert reclaimed.priority == 1 + lease["target_reservation"] = isolated_db.reserve_pipeline_target( + "owner/repo", "integration-review", 10, HEAD_A, reclaimed.id, 300, + task_started_at=reclaimed.started_at.astimezone(timezone.utc).isoformat(), + ownership_kind="headless") + + try: + route = pipeline_insights.execution_route( + reclaimed, "skill", str(first_dir), retain_budget=True) + + assert route["action"] == "prompt" + assert route["mode"] == "skill" + assert isolated_db.get_task(task.id).priority == 1 + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT pipeline_priority_restore FROM tasks WHERE id = ?", + (task.id,), + ).fetchone() + assert stored["pipeline_priority_restore"] == 2 + ledger = isolated_db.pipeline_github_budget_reservations("github-default") + assert ledger["count"] == 1 + assert ledger["items"][0]["route"] == "fallback_targeted" + finally: + pipeline_insights.release_execution_admission() + + isolated_db.mark_completed(reclaimed.id, "done") + worker._recur_after_run(reclaimed) + successor = isolated_db.get_task( + isolated_db.get_series(task.series_id)["next_task_id"]) + assert successor.priority == 2 + assert isolated_db.get_series(merge.series_id)["priority"] == 1 + + +def test_route_priority_inheritance_is_limited_to_integration_fallback(tmp_path): + review_dir = tmp_path / "review" + merge_dir = tmp_path / "merge" + review_dir.mkdir() + merge_dir.mkdir() + review_queue = { + "id": "review", "series_contains": "Project - REVIEW", "replicas": 2, + } + merge_queue = { + "id": "merge", "series_contains": "Project - MERGE", + "execution": {"stage": "merge"}, + } + profile = { + "queues": [review_queue, merge_queue], + "github_budget": { + "priority_one_headroom": { + "core": 1200, "search": 0, "graphql": 0, + }, + }, + } + review_series = _replica_series( + 1, review_dir, title="Project - REVIEW 1") + merge_series = _replica_series( + 2, merge_dir, title="Project - MERGE") + review_series["priority"] = 2 + merge_series["priority"] = 1 + series = [review_series, merge_series] + + inherit = pipeline_insights._post_preflight_admission_priority + assert inherit( + profile, series, "review", "fallback_targeted", + "integration-review", {"core": 1200}) == 1 + assert inherit( + profile, series, "review", "fallback_targeted", "review", + {"core": 1200}) is None + assert inherit( + profile, series, "review", "tool", "integration-review", + {"core": 1200}) is None + assert inherit( + profile, series, "merge", "fallback_targeted", + "integration-review", {"core": 1200}) is None + + # The effective vector is shared across all profiles using the same GitHub + # account; this profile need not own the strongest configured reserve. + profile["github_budget"]["priority_one_headroom"] = { + "core": 0, "search": 0, "graphql": 0, + } + assert inherit( + profile, series, "review", "fallback_targeted", + "integration-review", {"core": 1200}) == 1 + assert inherit( + profile, series, "review", "fallback_targeted", + "integration-review", {"core": 0}) is None + merge_series["priority"] = 2 + assert inherit( + profile, series, "review", "fallback_targeted", + "integration-review", {"core": 1200}) is None + + +def test_stale_integration_priority_is_restored_before_empty_replica_result( + isolated_db, monkeypatch, tmp_path): + first_dir = tmp_path / "review-1" + second_dir = tmp_path / "review-2" + first_dir.mkdir() + second_dir.mkdir() + first = isolated_db.create_task(TaskCreate( + prompt="Project - REVIEW 1", working_dir=str(first_dir), + recurrence="15m", priority=2)) + isolated_db.create_task(TaskCreate( + prompt="Project - REVIEW 2", working_dir=str(second_dir), + recurrence="15m", priority=2)) + task = isolated_db.get_next_runnable() + assert task.id == first.id + assert isolated_db.promote_running_attempt_priority( + task.id, task.started_at, 1) == 1 + task = isolated_db.get_task(task.id) + queue = { + "id": "review", "series_contains": "Project - REVIEW", "replicas": 2, + "execution": { + "mode": "auto", "stage": "review", + "command": ["pipelinectl", "next", "review"], + }, + } + profile = {"repository": "owner/repo", "queues": [queue]} + monkeypatch.setattr( + pipeline_insights, "_matching_queue", lambda _task: ("p", profile, queue)) + monkeypatch.setattr( + pipeline_insights, "_tool_available", lambda *_args: (True, "")) + monkeypatch.setattr( + pipeline_insights, "_tool_preflight", + lambda *_args, **_kwargs: {"action": "wait", "reason": "owner changed"}) + monkeypatch.setattr(pipeline_insights, "load_providers", lambda: {}) + + route = pipeline_insights.execution_route(task, "skill", str(first_dir)) + + assert route["action"] == "complete_empty" + assert isolated_db.get_task(task.id).priority == 2 + with isolated_db._connect() as conn: + stored = conn.execute( + "SELECT pipeline_priority_restore FROM tasks WHERE id = ?", + (task.id,), + ).fetchone() + assert stored["pipeline_priority_restore"] is None + + def test_worker_propagates_replica_mode_into_herdr_provider(monkeypatch, tmp_path): data_dir = tmp_path / "scheduler-data" started_at = datetime(2026, 9, 17, 12, 0, tzinfo=timezone.utc)