From 5bd33d4c91724c1bbf0a7177793bf9c784cce349 Mon Sep 17 00:00:00 2001 From: ibrog Date: Thu, 17 Sep 2026 15:22:21 +0300 Subject: [PATCH] fix(pipeline): defer merge behind local review Generated-with: Codex --- README.md | 22 ++++- promptpilot/db.py | 6 +- promptpilot/pipeline_insights.py | 115 +++++++++++++++++++++++- tests/test_schedule_series.py | 150 +++++++++++++++++++++++++++++++ 4 files changed, 285 insertions(+), 8 deletions(-) diff --git a/README.md b/README.md index d6274de..a71e10d 100644 --- a/README.md +++ b/README.md @@ -856,6 +856,7 @@ Worker публикует heartbeat в общей SQLite БД. Активный }, "dispatch_gate": { "skip_when_empty": true, + "defer_while_queues_active": ["review"], "defer_when_diagnostics_match": [{ "field": "review_candidates", "key": "stage", @@ -1039,15 +1040,28 @@ PromptPilot по умолчанию переиспользует результ запуском провайдера PromptPilot использует только полный общий снимок не старше пяти минут; устаревший или частичный снимок не может объявить очередь пустой; проектная процедура остаётся источником свежей проверки перед мутацией. При -`skip_when_empty` пустой запуск завершается как `ПУСТО`, не запуская LLM. Поле +`skip_when_empty` пустой запуск завершается как `ПУСТО`, не запуская LLM. +`defer_while_queues_active` задаёт локальную зависимость по id очередей того же +профиля. Она проверяется до GitHub admission и не зависит от диагностического +кэша: этап откладывается, пока зависимая серия выполняется, уже готова к запуску +или удерживается жёстким budget-defer. Обычный будущий запуск recurrence не +блокирует. Если одно завершение будит сразу зависимость и зависящий от неё этап, +PromptPilot сначала публикует wake зависимости, даже если в +`wake_after_success` id перечислены в обратном порядке. Это локальный порядок +диспетчеризации; выбор GitHub-цели, cleanup и fallback по-прежнему определяет +свежий проектный preflight. + `defer_when_diagnostics_nonempty` задаёт простую зависимость от непустого массива в JSON `health_check`. Когда один массив содержит несколько состояний, используйте `defer_when_diagnostics_match`: в примере MERGE возвращается в pending только для элементов со stage `integration-review`/`legacy-integration-review`, но проходит для `integration-merge-ready`. Это даёт автоматический порядок REVIEW → MERGE -даже при более высоком приоритете MERGE и не расходует токены на ожидание. Если -профиль или checker сломан, gate fail-open: задача запускается обычным способом, -чтобы ошибка наблюдаемости не остановила полезную работу. +даже при более высоком приоритете MERGE и не расходует токены на ожидание. +Диагностическая часть gate остаётся fail-open: при сломанном профиле/checker, +устаревшем или частичном снимке задача запускается обычным способом. Явно +настроенная локальная зависимость продолжает действовать и в этот момент, +поскольку её состояние берётся из той же SQLite-очереди, а не из наблюдаемости +GitHub. `wake_when` необязателен. Он указывает поле диагностики, которое означает, что этапу уже есть что делать. После полезного (`ГОТОВО`) запуска и при фоновом diff --git a/promptpilot/db.py b/promptpilot/db.py index 276c2c6..ef9e7fb 100644 --- a/promptpilot/db.py +++ b/promptpilot/db.py @@ -947,13 +947,14 @@ def list_series() -> list: GROUP BY series_id ) SELECT 'active' AS kind, t.series_id, t.id, t.machine, - t.scheduled_at, t.error + t.scheduled_at, t.next_run_at, t.error FROM tasks AS t INNER JOIN picked ON picked.active_id = t.id INNER JOIN task_series AS s ON s.id = t.series_id UNION ALL SELECT 'last' AS kind, t.series_id, t.id, t.machine, - NULL AS scheduled_at, NULL AS error + NULL AS scheduled_at, NULL AS next_run_at, + NULL AS error FROM tasks AS t INNER JOIN picked ON picked.last_id = t.id INNER JOIN task_series AS s ON s.id = t.series_id"""): @@ -993,6 +994,7 @@ def list_series() -> list: "next_task_id": active["id"] if active else None, "next_status": active["status"] if active else None, "next_run_at": active_detail["scheduled_at"] if active_detail else None, + "next_not_before": active_detail["next_run_at"] if active_detail else None, "next_error": active_detail["error"] if active_detail else None, "next_started_at": active["started_at"] if active else None, "last_task_id": last["id"] if last else None, diff --git a/promptpilot/pipeline_insights.py b/promptpilot/pipeline_insights.py index 2919ff6..b91bcb8 100644 --- a/promptpilot/pipeline_insights.py +++ b/promptpilot/pipeline_insights.py @@ -1853,11 +1853,16 @@ def dispatch_gate(task) -> dict | None: config = queue_config.get("dispatch_gate") if not marker or marker not in title or not isinstance(config, dict): continue + series = db.list_series() + dependency = _local_dependency_gate( + profile_id, profile, queue_config, config, series) + if dependency is not None: + return dependency # Dispatch only decides whether starting an agent is useful; every # mutation is still protected by the project's own fresh gate. # Reuse the five-minute snapshot so several due stages cannot each # spend hundreds of GitHub requests on the same queue state. - data = read_cached(profile_id, db.list_series()) + data = read_cached(profile_id, series) cache = data.get("cache") or {} # A stale/partial empty snapshot must never complete a live stage as # empty, and stale diagnostics must not defer it. The project-owned @@ -1909,6 +1914,79 @@ def dispatch_gate(task) -> dict | None: return None +def _local_dependency_gate(profile_id: str, profile: dict, queue: dict, + config: dict, series: list[dict], + *, now: datetime | None = None) -> dict | None: + """Defer one queue while an opted-in local predecessor still owns work. + + This check intentionally precedes and does not depend on the diagnostics + cache. It is only a scheduling guard: the project preflight remains the + authority for target selection, fallback and cleanup recovery. + """ + configured = _configured_local_dependencies(config) + if configured is None: + return None + queue_by_id = { + str(item.get("id")): item for item in profile.get("queues", []) + if isinstance(item, dict) and item.get("id") is not None + } + now = now or datetime.now(timezone.utc) + blockers = [] + for dependency_id in configured: + dependency_queue = queue_by_id.get(dependency_id) + if dependency_queue is None: + raise ValueError( + "dispatch_gate.defer_while_queues_active содержит неизвестную " + f"очередь: {dependency_id}") + dependency_series = _series_for_queue(dependency_queue, series) + state = _local_dependency_state(dependency_series, now) + if state is not None: + blockers.append((dependency_id, state)) + if not blockers: + return None + details = ", ".join(f"{queue_id} ({state})" for queue_id, state in blockers) + return { + "action": "defer", + "defer_for": config.get("defer_for", "10m"), + "reason": f"ожидание локальной очереди: {details}", + "profile_id": profile_id, + "queue_id": str(queue.get("id")), + } + + +def _configured_local_dependencies(config: object) -> list[str] | None: + if not isinstance(config, dict): + return None + configured = config.get("defer_while_queues_active") + if configured is None: + return None + if (not isinstance(configured, list) + or any(not isinstance(item, str) or not item.strip() + for item in configured)): + raise ValueError( + "dispatch_gate.defer_while_queues_active должен быть массивом id очередей") + return list(dict.fromkeys(item.strip() for item in configured)) + + +def _local_dependency_state(series: dict | None, + now: datetime) -> str | None: + """Classify only states that should keep a local successor deferred.""" + if not series or series.get("ended") or series.get("paused"): + return None + status = series.get("next_status") + if status == "running": + return "выполняется" + if status != "pending": + return None + not_before = _parse_time(series.get("next_not_before")) + if not_before is not None and not_before > now: + return "отложена бюджетом" + scheduled = _parse_time(series.get("next_run_at")) + if scheduled is None or scheduled <= now: + return "готова к запуску" + return None + + def _matching_queue(task) -> tuple[str, dict, dict] | None: """Return the profile and queue owning a recurring task, if configured.""" if not getattr(task, "series_id", None): @@ -2158,8 +2236,9 @@ def _wake_configured_successors(profile: dict, queue: dict, str(item.get("id")): item for item in profile.get("queues", []) if isinstance(item, dict) and item.get("id") is not None } + configured_ids = list(dict.fromkeys(item.strip() for item in configured)) woken = [] - for queue_id in dict.fromkeys(item.strip() for item in configured): + for queue_id in _ordered_successor_ids(configured_ids, queue_by_id): target_queue = queue_by_id.get(queue_id) if target_queue is None: raise ValueError( @@ -2182,6 +2261,38 @@ def _wake_configured_successors(profile: dict, queue: dict, return woken +def _ordered_successor_ids(configured: list[str], + queue_by_id: dict[str, dict]) -> list[str]: + """Publish local predecessors before successors made ready together.""" + configured_set = set(configured) + ordered = [] + visiting = set() + visited = set() + + def visit(queue_id: str) -> None: + if queue_id in visited: + return + if queue_id in visiting: + raise ValueError( + "dispatch_gate.defer_while_queues_active образует цикл " + f"для wake_after_success: {queue_id}") + visiting.add(queue_id) + target = queue_by_id.get(queue_id) + if target is not None: + dependencies = _configured_local_dependencies( + target.get("dispatch_gate")) or [] + for dependency_id in dependencies: + if dependency_id in configured_set: + visit(dependency_id) + visiting.remove(queue_id) + visited.add(queue_id) + ordered.append(queue_id) + + for queue_id in configured: + visit(queue_id) + return ordered + + def after_task_completed(task, verdict: str | None) -> list[str]: """Immediately advance ready pipeline stages after a productive run.""" if str(verdict or "").upper() != "ГОТОВО": diff --git a/tests/test_schedule_series.py b/tests/test_schedule_series.py index d1c0ed8..cc14664 100644 --- a/tests/test_schedule_series.py +++ b/tests/test_schedule_series.py @@ -852,6 +852,8 @@ def test_hard_defer_is_not_shortened_by_successor_wake(isolated_db): assert requested == {"accepted": True, "state": "latched_deferred"} assert deferred.scheduled_at == deadline assert deferred.next_run_at == deadline + assert isolated_db.get_series(task.series_id)["next_not_before"] == \ + deadline.isoformat() assert isolated_db.consume_pipeline_series_wake(task.series_id) is False assert isolated_db.get_setting( f"pipeline_series_wake_intent:v1:{task.series_id}") == "1" @@ -2378,6 +2380,154 @@ def test_dispatch_gate_defers_only_matching_diagnostic_stages(isolated_db, monke assert pipeline_insights.dispatch_gate(task) is None +def _local_dependency_profile(): + execution = { + "mode": "auto", "command": ["pipelinectl", "next", "{stage}"], + } + return { + "title": "Example", "repository": "owner/example", + "queues": [ + { + "id": "review", "title": "Review", "query": "is:pr", + "series_contains": "Example - REVIEW", "execution": execution, + "wake_after_success": ["merge"], + }, + { + "id": "merge", "title": "Merge", + "query": "is:pr label:ship", + "series_contains": "Example - MERGE", "execution": execution, + "wake_after_success": ["merge", "review"], + "dispatch_gate": { + "defer_while_queues_active": ["review"], + "defer_for": "7m", + }, + }, + ], + } + + +@pytest.mark.parametrize(("review_state", "reason"), [ + ("running", "выполняется"), + ("due", "готова к запуску"), + ("budget-deferred", "отложена бюджетом"), +]) +def test_local_dependency_precedes_invalidated_diagnostics_cache( + isolated_db, monkeypatch, review_state, reason): + review = isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h")) + if review_state in {"running", "budget-deferred"}: + claimed = isolated_db.get_next_runnable() + assert claimed.id == review.id + if review_state == "budget-deferred": + isolated_db.defer_task( + claimed.id, datetime.now(timezone.utc) + timedelta(hours=1), + "GitHub budget", hard_not_before=True) + merge = isolated_db.create_task(TaskCreate( + prompt="Example - MERGE", recurrence="4h")) + profile = _local_dependency_profile() + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "read_cached", + lambda *_args, **_kwargs: pytest.fail( + "local dependency must run before the invalidated diagnostics cache"), + ) + + gate = pipeline_insights.dispatch_gate(merge) + + assert gate["action"] == "defer" + assert gate["defer_for"] == "7m" + assert "review" in gate["reason"] + assert reason in gate["reason"] + + +def test_local_dependency_ignores_normal_future_recurrence( + isolated_db, monkeypatch): + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", + scheduled_at=datetime.now(timezone.utc) + timedelta(hours=4))) + merge = isolated_db.create_task(TaskCreate( + prompt="Example - MERGE", recurrence="4h")) + profile = _local_dependency_profile() + reads = [] + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "read_cached", lambda *_args, **_kwargs: ( + reads.append(True) or { + "queues": [{"id": "merge", "title": "Merge", "backlog": 1}], + "diagnostics": {"review_candidates": [{"number": 42}]}, + "cache": {"stale": True, "complete": False}, + })) + + assert pipeline_insights.dispatch_gate(merge) is None + assert reads == [True] + + +def test_local_dependency_defers_before_execution_admission( + isolated_db, monkeypatch): + isolated_db.create_task(TaskCreate( + prompt="Example - REVIEW", recurrence="4h", priority=5)) + merge = isolated_db.create_task(TaskCreate( + prompt="Example - MERGE", recurrence="4h", priority=1)) + claimed = isolated_db.get_next_runnable() + assert claimed.id == merge.id + profile = _local_dependency_profile() + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr( + pipeline_insights, "execution_route", + lambda *_args, **_kwargs: pytest.fail( + "local dependency must precede GitHub admission/reservation"), + ) + + worker._execute_task_inner(claimed) + + deferred = isolated_db.get_task(merge.id) + assert deferred.status.value == "pending" + assert deferred.retry_count == 0 + assert "локальной очереди" in deferred.error + + +def test_completion_invalidates_then_wakes_local_dependency_first( + monkeypatch): + profile = _local_dependency_profile() + series = [ + {"id": 7, "title": "Example - REVIEW", "paused": False, + "ended": False, "ended_at": None}, + {"id": 8, "title": "Example - MERGE", "paused": False, + "ended": False, "ended_at": None}, + ] + monkeypatch.setattr( + pipeline_insights, "_profiles", lambda: {"example": profile}) + monkeypatch.setattr(pipeline_insights.db, "is_paused", lambda: False) + monkeypatch.setattr(pipeline_insights.db, "list_series", lambda: series) + events = [] + monkeypatch.setattr( + pipeline_insights, "_discard_cache", + lambda profile_id: events.append(("invalidate", profile_id))) + monkeypatch.setattr( + pipeline_insights.db, "request_pipeline_series_wake", + lambda series_id: events.append(("wake", series_id)) or { + "accepted": True, "state": "scheduled", + }) + merge_task = SimpleNamespace( + series_id=8, series_title="Example - MERGE", prompt="Example - MERGE") + + assert pipeline_insights.after_task_completed( + merge_task, "ГОТОВО") == ["review", "merge"] + assert events == [ + ("invalidate", "example"), ("wake", 7), ("wake", 8), + ] + + events.clear() + review_task = SimpleNamespace( + series_id=7, series_title="Example - REVIEW", prompt="Example - REVIEW") + assert pipeline_insights.after_task_completed( + review_task, "ГОТОВО") == ["merge"] + assert events == [("invalidate", "example"), ("wake", 8)] + + def test_productive_completion_wakes_every_ready_stage(isolated_db, monkeypatch): profile = { "title": "Example", "repository": "owner/example",