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)