From 529192baec1206ca0fb8738ccb93686383be8745 Mon Sep 17 00:00:00 2001 From: Hermes Team Date: Sun, 30 Aug 2026 20:46:25 +0700 Subject: [PATCH] =?UTF-8?q?fix(roles):=2019=20=D0=B0=D0=B3=D0=B5=D0=BD?= =?UTF-8?q?=D1=82=D0=BE=D0=B2=20=D0=B2=D0=BC=D0=B5=D1=81=D1=82=D0=BE=2013,?= =?UTF-8?q?=20=D1=88=D0=B5=D1=81=D1=82=D1=8C=20=D0=BF=D0=B0=D1=80=20=D0=BD?= =?UTF-8?q?=D0=B5=D0=BE=D1=82=D0=BB=D0=B8=D1=87=D0=B8=D0=BC=D1=8B=20=D0=BF?= =?UTF-8?q?=D0=BE=20=D0=BD=D0=B0=D0=B7=D0=B2=D0=B0=D0=BD=D0=B8=D1=8E;=20?= =?UTF-8?q?=D0=BF=D0=BE=D0=B8=D1=81=D0=BA=20=D0=BB=D0=BE=D0=BA=D0=B0=D0=BB?= =?UTF-8?q?=D1=8C=D0=BD=D1=8B=D1=85=20=D1=81=D0=B5=D1=80=D0=B2=D0=B5=D1=80?= =?UTF-8?q?=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Проверка кандидата A28–A31 исполнением. 1. Дублирование ролей. RoleRegistry.migrate_legacy_roles написана верно, но НЕ ВЫЗЫВАЛАСЬ НИОТКУДА — проверено поиском по всему коду. Вместо неё работала «идемпотентная миграция» в router_config, дописывавшая недостающие умолчания и не убиравшая старые роли. На конфигурации владельца интерфейс показывал 19 агентов, причём шесть пар были неотличимы по названию: orchestrator и manager — оба «Менеджер проекта», reviewer и code-reviewer — оба «Ревьюер кода». Разложить аккаунты по такому списку невозможно. Миграция подключена. Внутри неё нашлась вторая ошибка: при обходе одним проходом пустая каноническая роль затирала цепочку, перенесённую из старой. У владельца сработало бы именно так — в manager попал бы codex-orch вместо выставленного им ag-orch-fallback. Старые роли теперь обрабатываются первыми: в них и лежит настроенный порядок аккаунтов. Проверено на живой конфигурации владельца: было 19 ролей, стало 13, все цепочки совпадают с исходными. 2. Сохранённый workflow мигрируется вместе с ролями. Идентификаторы агентов повторяют идентификаторы ролей, и без переименования рёбра ссылались бы на исчезнувших агентов — «Ребро ссылается на отсутствующего агента». 3. Поиск локальных серверов моделей. Раньше адрес вводился руками. Теперь опрашиваются известные порты на петле — Ollama 11434, LM Studio 1234, llama.cpp 8080-8082, vLLM 8000, Jan, GPT4All, Text Generation WebUI, — параллельно, девять портов за 1.5 с. Наружу идёт только то, что ответило; список моделей берётся у сервера. Порт, занятый чужим сервисом, показывается с причиной, закрытые не показываются вовсе. Действие discover_local_models. Роль в фикстурах test_workflow_service_a30 переименована в test-developer: «developer» — псевдоним канонической developer-1, и тест проверял бы работу псевдонимов вместо механики workflow. 475 passed, 2 skipped; ruff чисто. Co-Authored-By: Claude Opus 5 --- .../router/action_handler.py | 28 ++++ .../router/local_discovery.py | 124 ++++++++++++++++++ .../router/role_registry.py | 20 ++- .../router/router_config.py | 20 +++ .../router/workflow_service.py | 41 +++++- tests/test_workflow_service_a30.py | 38 +++--- 6 files changed, 246 insertions(+), 25 deletions(-) create mode 100644 src/antigravity_provider/router/local_discovery.py diff --git a/src/antigravity_provider/router/action_handler.py b/src/antigravity_provider/router/action_handler.py index f68fd23..5878d02 100644 --- a/src/antigravity_provider/router/action_handler.py +++ b/src/antigravity_provider/router/action_handler.py @@ -596,6 +596,34 @@ class ActionExecutor: ok, msg = do_save_settings(data) return {'ok': ok, 'message': msg} + elif action == 'discover_local_models': + # Поиск уже запущенных локальных серверов на машине. + # + # Раньше адрес вводился руками, и владелец должен был помнить, на + # каком порту у него llama.cpp, Ollama или LM Studio. Опрашиваются + # известные порты на петле; наружу идёт только то, что ответило. + # Ни одной выдуманной модели: список берётся у самого сервера. + from antigravity_provider.router.local_discovery import discover_local_servers + + host = (data.get('host') or '127.0.0.1').strip() or '127.0.0.1' + try: + servers = discover_local_servers(host=host) + except Exception as exc: + return {'ok': False, 'message': f'Поиск локальных серверов не удался: {exc}'} + + usable = [s for s in servers if not s.error] + if not servers: + return { + 'ok': True, + 'message': 'Локальных серверов не найдено. Запустите Ollama, LM Studio или llama.cpp либо введите адрес вручную.', + 'data': {'servers': [], 'host': host}, + } + return { + 'ok': True, + 'message': f'Найдено серверов: {len(usable)}', + 'data': {'servers': [s.to_dict() for s in servers], 'host': host}, + } + elif action == 'check_updates': # Проверка выполняется СИНХРОННО и всегда возвращает данные. # diff --git a/src/antigravity_provider/router/local_discovery.py b/src/antigravity_provider/router/local_discovery.py new file mode 100644 index 0000000..018fd45 --- /dev/null +++ b/src/antigravity_provider/router/local_discovery.py @@ -0,0 +1,124 @@ +"""Поиск локальных серверов моделей, уже запущенных на машине. + +Раньше адрес локального сервера вводился руками: владелец должен был помнить, +на каком порту у него llama.cpp, Ollama или LM Studio. Здесь опрашиваются +известные порты на петле, и наружу отдаётся только то, что **ответило**. + +Правило честности действует и тут: ни одного порта, который не ответил, в +результат не попадает, и ни одна модель не выдумывается — список берётся у +самого сервера. +""" + +from __future__ import annotations + +import concurrent.futures +import json +import logging +import urllib.error +import urllib.request +from dataclasses import dataclass, field +from typing import Any, Dict, List, Optional + +logger = logging.getLogger("hermes.router.local_discovery") + +# Порты, на которых эти серверы работают из коробки. Список — умолчание, а не +# исчерпывающая истина: владелец всегда может ввести адрес сам. +WELL_KNOWN_ENDPOINTS: List[Dict[str, Any]] = [ + {"name": "Ollama", "port": 11434, "base_path": "/v1"}, + {"name": "LM Studio", "port": 1234, "base_path": "/v1"}, + {"name": "llama.cpp", "port": 8080, "base_path": "/v1"}, + {"name": "llama.cpp", "port": 8081, "base_path": "/v1"}, + {"name": "llama.cpp", "port": 8082, "base_path": "/v1"}, + {"name": "vLLM", "port": 8000, "base_path": "/v1"}, + {"name": "Jan", "port": 1337, "base_path": "/v1"}, + {"name": "GPT4All", "port": 4891, "base_path": "/v1"}, + {"name": "Text Generation WebUI", "port": 5000, "base_path": "/v1"}, +] + +DEFAULT_HOST = "127.0.0.1" +PROBE_TIMEOUT_SEC = 1.5 + + +@dataclass +class LocalServer: + """Обнаруженный локальный сервер моделей.""" + + name: str + base_url: str + models: List[str] = field(default_factory=list) + error: Optional[str] = None + + def to_dict(self) -> Dict[str, Any]: + return { + "name": self.name, + "base_url": self.base_url, + "models": list(self.models), + "model_count": len(self.models), + "error": self.error, + } + + +def _fetch_models(base_url: str, timeout: float) -> List[str]: + """Спросить у сервера список моделей по OpenAI-совместимому /models.""" + url = base_url.rstrip("/") + "/models" + req = urllib.request.Request(url, headers={"Accept": "application/json"}) + with urllib.request.urlopen(req, timeout=timeout) as resp: + if resp.status != 200: + raise ValueError(f"HTTP {resp.status}") + payload = json.loads(resp.read().decode("utf-8", errors="replace")) + + models: List[str] = [] + for item in payload.get("data") or []: + if isinstance(item, dict) and item.get("id"): + models.append(str(item["id"])) + elif isinstance(item, str): + models.append(item) + return models + + +def probe_endpoint(host: str, port: int, base_path: str, name: str, + timeout: float = PROBE_TIMEOUT_SEC) -> Optional[LocalServer]: + """Опросить один адрес. Не ответил — вернуть None, а не пустую заглушку.""" + base_url = f"http://{host}:{port}{base_path}" + try: + models = _fetch_models(base_url, timeout) + except (urllib.error.URLError, OSError, TimeoutError): + # Порт закрыт или сервер не отвечает — это не находка и не ошибка. + return None + except (ValueError, json.JSONDecodeError) as exc: + # Порт занят чем-то, что говорит не на том языке. Показываем как + # найденный, но с причиной: иначе владелец будет гадать, почему + # заведомо работающий сервер не виден. + logger.debug("Порт %s занят не тем сервисом: %s", port, exc) + return LocalServer(name=f"{name} (порт {port})", base_url=base_url, + error="Сервер на этом порту ответил не по OpenAI-совместимому протоколу") + return LocalServer(name=f"{name} (порт {port})", base_url=base_url, models=models) + + +def discover_local_servers(host: str = DEFAULT_HOST, + endpoints: Optional[List[Dict[str, Any]]] = None, + timeout: float = PROBE_TIMEOUT_SEC) -> List[LocalServer]: + """Опросить известные порты параллельно и вернуть только ответившие. + + Опрос идёт параллельно: последовательно девять закрытых портов дали бы + заметную паузу в интерфейсе. + """ + targets = endpoints if endpoints is not None else WELL_KNOWN_ENDPOINTS + found: List[LocalServer] = [] + + with concurrent.futures.ThreadPoolExecutor(max_workers=min(10, len(targets) or 1)) as pool: + futures = { + pool.submit(probe_endpoint, host, t["port"], t.get("base_path", "/v1"), t["name"], timeout): t + for t in targets + } + for fut in concurrent.futures.as_completed(futures): + try: + result = fut.result() + except Exception as exc: # опрос не должен ронять вызывающего + logger.warning("Ошибка опроса локального сервера: %s", exc) + continue + if result is not None: + found.append(result) + + found.sort(key=lambda s: (s.error is not None, -len(s.models), s.base_url)) + return found diff --git a/src/antigravity_provider/router/role_registry.py b/src/antigravity_provider/router/role_registry.py index fb749b6..e60b496 100644 --- a/src/antigravity_provider/router/role_registry.py +++ b/src/antigravity_provider/router/role_registry.py @@ -333,7 +333,21 @@ class RoleRegistry: was_modified = False default_policies = cls.get_default_role_policies() - for rname, rpol in current_roles.items(): + # Сначала роли под старыми именами, потом под каноническими. + # + # Порядок важен, когда в конфигурации есть и то и другое. У владельца + # ровно такой случай: прошлая миграция дописала manager, developer-1 и + # остальные рядом со старыми orchestrator и coder-primary, но работал + # он всё это время со старыми — там и лежит его настроенный порядок + # аккаунтов. При обходе одним проходом канонический пустой manager + # затирал бы цепочку из orchestrator, и владелец получил бы на первом + # месте аккаунт, который сам туда не ставил. + legacy_first = sorted( + current_roles.items(), + key=lambda kv: cls.resolve_canonical_role(kv[0]) == kv[0], + ) + + for rname, rpol in legacy_first: canonical_id = cls.resolve_canonical_role(rname) if canonical_id != rname: was_modified = True @@ -346,8 +360,10 @@ class RoleRegistry: session_affinity_enabled=rpol.session_affinity_enabled, default_model=rpol.default_model or default_policies[canonical_id].default_model, ) - else: + elif rname not in migrated: migrated[rname] = rpol + else: + was_modified = True for canon_id, def_policy in default_policies.items(): if canon_id not in migrated: diff --git a/src/antigravity_provider/router/router_config.py b/src/antigravity_provider/router/router_config.py index 465bfe0..bd48408 100644 --- a/src/antigravity_provider/router/router_config.py +++ b/src/antigravity_provider/router/router_config.py @@ -377,6 +377,26 @@ def load_router_config(config_path: Optional[Path] = None) -> RouterConfig: if not roles: roles = default_cfg.roles else: + # Сначала переименование старых ролей в канонические, и только + # потом дополнение недостающими. + # + # Раньше здесь просто дописывались отсутствующие умолчания, а + # RoleRegistry.migrate_legacy_roles не вызывалась ниоткуда — + # проверено поиском по всему коду. В результате старые роли + # оставались рядом с новыми: на конфигурации владельца интерфейс + # показывал 19 агентов вместо 13, причём шесть пар были неотличимы + # по названию (orchestrator и manager — оба «Менеджер проекта», + # reviewer и code-reviewer — оба «Ревьюер кода»). Разложить + # аккаунты по такому списку невозможно. + # + # migrate_legacy_roles переносит preferred_chain дословно, поэтому + # порядок аккаунтов, выставленный владельцем, сохраняется. + renamed, renamed_any = RoleRegistry.migrate_legacy_roles(roles) + if renamed_any: + new_roles_added.extend(sorted(set(renamed) - set(roles))) + roles = renamed + migration_needed = True + for def_rname, def_rpolicy in default_cfg.roles.items(): if def_rname not in roles: roles[def_rname] = def_rpolicy diff --git a/src/antigravity_provider/router/workflow_service.py b/src/antigravity_provider/router/workflow_service.py index 25deca1..99059d9 100644 --- a/src/antigravity_provider/router/workflow_service.py +++ b/src/antigravity_provider/router/workflow_service.py @@ -206,13 +206,42 @@ class WorkflowService: return try: raw = json.loads(self.state_path.read_text(encoding="utf-8")) - self.agents = { - item["id"]: AgentDefinition(**item) - for item in raw.get("agents", []) - if isinstance(item, dict) and item.get("id") - } + + # Идентификаторы агентов повторяют идентификаторы ролей, а роли со + # старыми именами переименовываются в канонические при загрузке + # конфигурации. Без такого же переименования здесь сохранённый + # workflow ссылался бы на исчезнувших агентов, и граф падал бы с + # «Ребро ссылается на отсутствующего агента». + from antigravity_provider.router.role_registry import RoleRegistry + + def _canon(agent_id: str) -> str: + try: + return RoleRegistry.resolve_canonical_role(agent_id) + except Exception: + return agent_id + + self.agents = {} + for item in raw.get("agents", []): + if not isinstance(item, dict) or not item.get("id"): + continue + item = dict(item) + item["id"] = _canon(item["id"]) + if item.get("role"): + item["role"] = _canon(item["role"]) + # Первым выигрывает агент под старым именем: именно им + # пользовался владелец, канонический мог быть дописан пустым. + self.agents.setdefault(item["id"], AgentDefinition(**item)) + wf = raw.get("workflow") or {} - edges = [WorkflowEdge(**edge) for edge in wf.pop("edges", []) if isinstance(edge, dict)] + edges = [] + for edge in wf.pop("edges", []): + if not isinstance(edge, dict): + continue + edge = dict(edge) + for key in ("source", "target", "from_agent", "to_agent"): + if edge.get(key): + edge[key] = _canon(edge[key]) + edges.append(WorkflowEdge(**edge)) self.workflow = WorkflowDefinition(edges=edges, **wf) self.events = [WorkflowEvent(**event) for event in raw.get("events", [])[-200:]] self.run = raw.get("run") or self._idle_run() diff --git a/tests/test_workflow_service_a30.py b/tests/test_workflow_service_a30.py index 8948b35..f740fd5 100644 --- a/tests/test_workflow_service_a30.py +++ b/tests/test_workflow_service_a30.py @@ -13,6 +13,10 @@ from antigravity_provider.router.router_config import ( from antigravity_provider.router.workflow_service import WorkflowService +# Роль в фикстурах названа test-developer намеренно: "developer" — это +# псевдоним канонической роли developer-1, и при загрузке конфигурации он +# переименовывается. Тест проверяет механику workflow, а не работу псевдонимов, +# поэтому имя взято такое же нейтральное, как у соседнего test-reviewer. @pytest.fixture def workflow_service(tmp_path, monkeypatch): monkeypatch.setenv("HERMES_HOME", str(tmp_path)) @@ -24,7 +28,7 @@ def workflow_service(tmp_path, monkeypatch): preferred_models=["model-real-from-config"], ) }, - roles={"developer": RolePolicy(role_name="developer", preferred_chain=["account-a"])}, + roles={"test-developer": RolePolicy(role_name="test-developer", preferred_chain=["account-a"])}, ) save_router_config(config) return WorkflowService(tmp_path / "workflow_state.json") @@ -32,11 +36,11 @@ def workflow_service(tmp_path, monkeypatch): def test_router_roles_migrate_to_agents_and_create_real_files(workflow_service, tmp_path): snapshot = workflow_service.snapshot() - assert "developer" in [agent["id"] for agent in snapshot["agents"]] - agent = next(item for item in snapshot["agents"] if item["id"] == "developer") - assert agent["agent_file"] == "agents/developer.md" + assert "test-developer" in [agent["id"] for agent in snapshot["agents"]] + agent = next(item for item in snapshot["agents"] if item["id"] == "test-developer") + assert agent["agent_file"] == "agents/test-developer.md" assert agent["agent_file_exists"] is True - assert (tmp_path / "agents" / "developer.md").is_file() + assert (tmp_path / "agents" / "test-developer.md").is_file() assert agent["execution_config"]["account"] == "account-a" assert agent["execution_config"]["model"] == "model-real-from-config" @@ -62,9 +66,9 @@ def test_create_update_file_and_restart_persistence(workflow_service, tmp_path): def test_delete_requires_explicit_confirmation_when_referenced(workflow_service): workflow_service.create_agent({"name": "Test Reviewer", "role": "test-reviewer", "account": "account-a"}) workflow_service.save_workflow({ - "start_agent_id": "developer", + "start_agent_id": "test-developer", "max_iterations": 3, - "edges": [{"id": "review", "source": "developer", "target": "test-reviewer", "condition": "SUCCESS"}], + "edges": [{"id": "review", "source": "test-developer", "target": "test-reviewer", "condition": "SUCCESS"}], }) warning = workflow_service.delete_agent("test-reviewer") assert warning["confirmation_required"] is True @@ -77,11 +81,11 @@ def test_delete_requires_explicit_confirmation_when_referenced(workflow_service) def test_cycles_are_valid_and_iteration_limit_is_persisted(workflow_service): workflow_service.create_agent({"name": "Test Reviewer", "role": "test-reviewer", "account": "account-a"}) definition = workflow_service.save_workflow({ - "start_agent_id": "developer", + "start_agent_id": "test-developer", "max_iterations": 2, "edges": [ - {"source": "developer", "target": "test-reviewer", "condition": "SUCCESS"}, - {"source": "test-reviewer", "target": "developer", "condition": "REVIEW_FAILED"}, + {"source": "test-developer", "target": "test-reviewer", "condition": "SUCCESS"}, + {"source": "test-reviewer", "target": "test-developer", "condition": "REVIEW_FAILED"}, ], }) assert definition.max_iterations == 2 @@ -92,13 +96,13 @@ def test_cycles_are_valid_and_iteration_limit_is_persisted(workflow_service): def test_invalid_edge_and_unknown_model_are_rejected(workflow_service): with pytest.raises(ValueError, match="отсутствующего агента"): - workflow_service.save_workflow({"edges": [{"source": "developer", "target": "missing"}]}) + workflow_service.save_workflow({"edges": [{"source": "test-developer", "target": "missing"}]}) with pytest.raises(ValueError, match="не доступна"): - workflow_service.update_agent("developer", {"account": "account-a", "model": "invented-model"}) + workflow_service.update_agent("test-developer", {"account": "account-a", "model": "invented-model"}) def test_interrupted_run_is_reported_not_silently_completed(workflow_service, tmp_path): - workflow_service.run.update({"id": "run-1", "status": "running", "current_agent_id": "developer"}) + workflow_service.run.update({"id": "run-1", "status": "running", "current_agent_id": "test-developer"}) workflow_service._save() restarted = WorkflowService(tmp_path / "workflow_state.json") assert restarted.run["status"] == "interrupted" @@ -109,11 +113,11 @@ def test_interrupted_run_is_reported_not_silently_completed(workflow_service, tm def test_live_cycle_stops_with_explicit_iteration_limit_event(workflow_service, monkeypatch): workflow_service.create_agent({"name": "Loop Reviewer", "role": "loop-reviewer", "account": "account-a"}) workflow_service.save_workflow({ - "start_agent_id": "developer", + "start_agent_id": "test-developer", "max_iterations": 2, "edges": [ - {"source": "developer", "target": "loop-reviewer", "condition": "SUCCESS"}, - {"source": "loop-reviewer", "target": "developer", "condition": "REVIEW_FAILED"}, + {"source": "test-developer", "target": "loop-reviewer", "condition": "SUCCESS"}, + {"source": "loop-reviewer", "target": "test-developer", "condition": "REVIEW_FAILED"}, ], }) @@ -155,7 +159,7 @@ def test_provider_error_text_reaches_run_and_events(workflow_service, monkeypatc return {"choices": [{"message": {"content": f"ERROR\n{provider_text}"}}]} monkeypatch.setattr("antigravity_provider.router.router_engine.get_router_engine", lambda: ErrorEngine()) - workflow_service.workflow.start_agent_id = "developer" + workflow_service.workflow.start_agent_id = "test-developer" workflow_service.workflow.edges = [] workflow_service.start("Проверить ошибку") thread = workflow_service._thread