Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 18 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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` необязателен. Он указывает поле диагностики, которое означает, что
этапу уже есть что делать. После полезного (`ГОТОВО`) запуска и при фоновом
Expand Down
6 changes: 4 additions & 2 deletions promptpilot/db.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"""):
Expand Down Expand Up @@ -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,
Expand Down
115 changes: 113 additions & 2 deletions promptpilot/pipeline_insights.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down Expand Up @@ -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(
Expand All @@ -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() != "ГОТОВО":
Expand Down
150 changes: 150 additions & 0 deletions tests/test_schedule_series.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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",
Expand Down
Loading