Skip to content
Merged
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
87 changes: 81 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -682,15 +682,17 @@ Web UI: **⚙ Providers** → «Изменить» у нужного прова
атомарно сохраняет полный последний успешный ответ и исторический снимок по
`profile_id`. Обычное открытие читает только SQLite и не обращается к GitHub —
в том числе после перезапуска server или bot. Возраст, источник и признак
устаревания снимка видны в интерфейсе. Узкое место — этап с максимальным ETA: учитываются
`backlog / capacity`, интервал серии и средняя длительность её выполнения.
устаревания снимка видны в интерфейсе. Узкое место — этап с максимальным ETA:
учитываются `backlog / (capacity × активные реплики)`, интервал серии и средняя
длительность её выполнения. Без явного `replicas` число реплик равно единице.
PromptPilot не поставляет профили конкретных проектов. Метки, запросы и лимиты
задаются пользователем в `~/.promptpilot/pipeline_profiles.json` (или файле из
`PP_PIPELINE_PROFILES`), поэтому чужие репозитории не появляются в новой установке.
Если в профиле включён `priority_control`, дашборд показывает элементы очередей
и позволяет поставить P0…P3 либо вернуть автоматическую оценку. Кнопка
«⚡ Следующим» ставит P0 и переносит будущий запуск соответствующей серии на
ближайшее время; поставленная на паузу серия при этом остаётся на паузе.
«⚡ Следующим» ставит P0 и переносит будущий запуск всех серий-реплик
соответствующей очереди на ближайшее время; поставленная на паузу серия при
этом остаётся на паузе.
Одинаковые диагностические сигналы в окне конвейера сворачиваются в одну строку
со списком и количеством PR; владелец single-flight барьера всегда показывается
первым, чтобы было видно, какой PR сейчас должен пройти REVIEW → MERGE. В основной
Expand Down Expand Up @@ -736,8 +738,9 @@ GitHub CLI и Go в `C:\Program Files`. Поэтому tray/служба не з
останавливает дерево. Поэтому штатный stop не оставляет скрытый worker, который
продолжает расходовать токены после «остановки» старого launcher PID.

Рекомендация считается прозрачно: `backlog / capacity` даёт число необходимых
прогонов, а один цикл равен `интервал + средняя длительность прогона` — повтор
Рекомендация считается прозрачно: `backlog / (capacity × активные реплики)`
даёт число необходимых параллельных прогонов, а один цикл равен
`интервал + средняя длительность прогона` — повтор
назначается после завершения предыдущего запуска. Поэтому серия `15m` со средним
выполнением 30 минут даёт примерно один результат за 45 минут, а не четыре в час.
Профиль задаёт желаемое время очистки (`target_clear_hours`, обычно 8 часов),
Expand Down Expand Up @@ -901,6 +904,75 @@ claim одной задачи происходят в одной SQLite-тран
получат одно вхождение. Без `scheduler` остаётся прежний порядок
`priority → created_at → id`.

#### Параллельные реплики REVIEW

Очередь REVIEW может явно владеть несколькими независимыми сериями. Для этого
используется отдельное поле `replicas`; `capacity` по-прежнему означает, сколько
целей обрабатывает **один** запуск, и не подменяет число исполнителей:

```json
{
"id": "review",
"capacity": 1,
"replicas": 2,
"series_contains": "MyProject - REVIEW",
"execution": {
"mode": "auto",
"stage": "review",
"command": [
"{python}", "-m", "promptpilot.project_pipeline",
"--config", "pipelinectl.json", "next", "{stage}"
],
"required_paths": ["pipelinectl.json"]
}
}
```

Нужно создать ровно две активные повторяющиеся серии, например
`MyProject - REVIEW 1` и `MyProject - REVIEW 2`. Каждой задаётся существующий
абсолютный и отличный от остальных `working_dir` — обычно два git worktree.
Отличие проверяется по identity каталога (`device/inode`), поэтому symlink и
варианты регистра одного APFS/NTFS-каталога не считаются разными репликами.
Реплики выполняются на той же машине и используют общую локальную SQLite:
удалённые `machine`, динамический `worktree` задачи, отсутствующие каталоги и
дубли путей блокируются до запуска провайдера. Если настроен `scheduler`, REVIEW
должен присутствовать минимум в двух lanes; без scheduler достаточно общей
конкурентности worker не меньше двух.

В `pipelinectl.json` для параллельного REVIEW обязательны
`review_completion_gate: "target-v1"` и `fallback_handoff: "target-v1"`.
Необязательный `target_reservation_ttl_seconds` задаёт TTL локальной резервации
(по умолчанию равен `review_lease_seconds`, допустимо 300–28800 секунд):

```json
{
"review_completion_gate": "target-v1",
"fallback_handoff": "target-v1",
"review_lease_seconds": 7200,
"target_reservation_ttl_seconds": 7200
}
```

До запуска модели `next review` атомарно резервирует точный
`repository/stage/PR/HEAD` за конкретной попыткой `PP_TASK_ID`; mutex действует
на весь PR даже при смене stage или HEAD. Если первая цель уже занята другой
репликой, election берёт следующую, не начиная повторное ревью. Обычный
`complete review` и `gate-fallback` проверяют и продлевают ту же резервацию
непосредственно перед первой GitHub-мутацией. Обычный terminal exit освобождает
резервацию. TTL задаёт частоту lease heartbeat, но истёкшая строка намеренно не
крадётся, пока её exact task attempt всё ещё `running`: после жёсткого падения
startup recovery сначала по сохранённому ownership descriptor подтверждает
остановку headless/Herdr-провайдера, затем атомарно requeue/cancel задачи и
освобождает fence. Если остановку доказать нельзя, задача и PR остаются в
quarantine; вручную удалять такую резервацию небезопасно. Резервация — только
локальный scheduling fence и не заменяет GraphQL/HEAD/timeline/CAS-гейты.

`run_now`, `wake_when`, `wake_after_success` и `adaptive_cadence` применяются ко
всем совпавшим репликам. Дашборд показывает каждую серию, её каталог и состояние,
а ETA использует суммарную параллельную ёмкость. Сейчас `replicas > 1`
намеренно разрешено только для REVIEW; MERGE и прочие изменяющие этапы остаются
single-flight.

Полоса выбирает только серию TRIAGE. Порядок целей *внутри* неё — recovery,
ручной P0, затем обычный backlog — остаётся контрактом репозиторного `next` и
его свежих gate-проверок. Так машинное расписание не может обойти recovery или
Expand Down Expand Up @@ -1116,6 +1188,9 @@ HEAD в `review_candidates`. Для двухполосной схемы он т
дашборде и боте для каждого этапа. Это позволяет заметить незапланированный
fallback сразу, а не по суточному расходу токенов. Пустой REVIEW/MERGE
закрывается самим preflight с явной записью «провайдер не запускался».
Для `replicas > 1` небезопасный нетаргетированный `auto → skill` отключён:
недоступный CLI, malformed lease или fallback без signed target блокирует запуск
до модели, а доказанный `target-v1` handoff остаётся разрешён.

Интеграционный владелец на `integration-review` или
`legacy-integration-review` всегда оставляет MERGE на полном skill-маршруте.
Expand Down
11 changes: 11 additions & 0 deletions promptpilot/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -338,6 +338,9 @@ def api_pipeline_item_priority(profile_id: str, queue_id: str, kind: str, number

@app.delete("/api/tasks/{task_id}", response_model=dict)
def api_delete_task(task_id: int):
task = db.get_task(task_id)
if task and task.status.value == "running":
raise HTTPException(409, "Cancel the running task and wait for it to stop first")
if not db.delete_task(task_id):
raise HTTPException(404, "Task not found")
return {"ok": True}
Expand All @@ -346,6 +349,14 @@ def api_delete_task(task_id: int):
@app.post("/api/tasks/{task_id}/reset", response_model=dict)
def api_reset_task(task_id: int):
if not db.reset_task(task_id):
task = db.get_task(task_id)
if (task and task.status.value == "running"
and db.task_has_live_pipeline_target_reservation(task_id)):
raise HTTPException(
409,
"This pipeline task still owns a live provider/target; cancel it "
"and wait for cleanup before resetting",
)
raise HTTPException(400, "Task not found or not in running state")
return {"ok": True}

Expand Down
50 changes: 43 additions & 7 deletions promptpilot/bot.py
Original file line number Diff line number Diff line change
Expand Up @@ -199,8 +199,11 @@ def _task_detail_keyboard(task) -> InlineKeyboardMarkup:
rows.append(action_row)
if task.result and len(task.result) > 800:
rows.append([InlineKeyboardButton("📄 Полный вывод", callback_data=f"full_result:{task.id}")])
rows.append([InlineKeyboardButton("← К списку", callback_data="tasklist"),
InlineKeyboardButton("🗑 Удалить", callback_data=f"delete_task:{task.id}")])
final_row = [InlineKeyboardButton("← К списку", callback_data="tasklist")]
if status != "running":
final_row.append(InlineKeyboardButton(
"🗑 Удалить", callback_data=f"delete_task:{task.id}"))
rows.append(final_row)
return InlineKeyboardMarkup(rows)


Expand Down Expand Up @@ -657,7 +660,13 @@ async def cb_reset_task(update: Update, context: ContextTypes.DEFAULT_TYPE):
await query.edit_message_text(f"Задача #{task_id} возвращена в очередь.",
reply_markup=_after_action_keyboard(task_id))
else:
await query.answer("Задача не в статусе running.", show_alert=True)
if db.task_has_live_pipeline_target_reservation(task_id):
await query.answer(
"Сначала останови задачу и дождись завершения очистки провайдера.",
show_alert=True,
)
else:
await query.answer("Задача не в статусе running.", show_alert=True)


@require_auth
Expand All @@ -666,8 +675,13 @@ async def cb_delete_task(update: Update, context: ContextTypes.DEFAULT_TYPE):
the button sits next to the frequently-used ones and a stray tap would
destroy the task together with its result."""
query = update.callback_query
await query.answer()
task_id = int(query.data.split(":")[1])
task = db.get_task(task_id)
if task and task.status.value == "running":
await query.answer(
"Сначала отмените задачу и дождитесь её остановки.", show_alert=True)
return
await query.answer()
await query.edit_message_reply_markup(reply_markup=InlineKeyboardMarkup([[
InlineKeyboardButton("🗑 Точно удалить", callback_data=f"del_yes:{task_id}"),
InlineKeyboardButton("↩ Отмена", callback_data=f"task:{task_id}"),
Expand All @@ -685,7 +699,11 @@ async def cb_delete_task_confirm(update: Update, context: ContextTypes.DEFAULT_T
reply_markup=InlineKeyboardMarkup([[
InlineKeyboardButton("← К списку", callback_data="tasklist")]]))
else:
await query.answer("Задача не найдена.", show_alert=True)
task = db.get_task(task_id)
reason = ("Сначала отмените задачу и дождитесь её остановки."
if task and task.status.value == "running"
else "Задача не найдена.")
await query.answer(reason, show_alert=True)


@require_auth
Expand Down Expand Up @@ -986,12 +1004,30 @@ def _pipeline_text(data: dict) -> str:
cadence = queue.get("adaptive_cadence") or {}
cadence_text = (f"; adaptive {cadence.get('mode')}"
if cadence else "")
replica_count = int(queue.get("replica_count") or 1)
replicas_active = int(queue.get("replicas_active") or 0)
parallel_capacity = queue.get("parallel_capacity")
parallel_capacity = int(queue["capacity"] if parallel_capacity is None
else parallel_capacity)
capacity_text = (
f"{queue['capacity']} × {replicas_active} = {parallel_capacity}")
replica_status = queue.get("replica_status") or {}
replica_details = ""
if replica_count > 1:
issues = replica_status.get("issues") or []
replica_details = (
f"\n Реплики: {replicas_active}/{replica_count} активны, "
f"найдено {replica_status.get('present', 0)}"
+ (f"; ошибка: {'; '.join(issues)}" if issues else "")
)
lines.append(f"{marker} {queue['title']}: {backlog if backlog is not None else '—'} / "
f"{queue['capacity']} за прогон = {runs_needed if runs_needed is not None else '—'} прогонов; "
f"{capacity_text} за параллельный прогон = "
f"{runs_needed if runs_needed is not None else '—'} прогонов; "
f"сейчас {queue['interval'] or 'не настроено'}{cadence_text}; "
f"средний запуск {duration_text}; ETA {eta if eta is not None else '—'} ч\n"
f" Маршрут: {route_text}\n"
f" Рекомендация: {queue['recommendation']}")
f" Рекомендация: {queue['recommendation']}"
f"{replica_details}")
runs = recent.get("runs", {})
lines.extend(["", f"Прогоны за окно: {runs.get('runs', 0)}; готово {runs.get('ready', 0)}, "
f"нужен человек {runs.get('human', 0)}, не смог {runs.get('unable', 0)}, "
Expand Down
Loading
Loading