diff --git a/env_example b/env_example index e3e2904..ac95069 100644 --- a/env_example +++ b/env_example @@ -40,6 +40,8 @@ SERVERUS_CA_DIR= SECRET_KEY= # Для админки через другой адрес, чем 127.0.0.1:8000, например https://admin.example.com CSRF_TRUSTED_ORIGINS= + +# Уровень логов сайта и воркера: DEBUG покажет и каждую проверку супервизора LOG_LEVEL=INFO # Бэкапы базы (make dj_backup), по умолчанию extra/backups diff --git a/src/App/services.py b/src/App/services.py index 8b1405d..884afed 100644 --- a/src/App/services.py +++ b/src/App/services.py @@ -47,11 +47,19 @@ def release_point(point: Point) -> tuple[dict[str, str], dict[str, str]]: """ passwords, errors = {}, {} for server in point.servers.all(): + if not server.private_key: + # возвращен прошлой попыткой: ключа на сервере больше нет, пароль уже новый + passwords[server.address] = server.root_password + continue password = serverus.generate_password() result = serverus.Server(server.host()).release(password) if result.ok: + # release убирает ключ и возвращает ssh на 22 порт server.root_password = password - server.save(update_fields=['root_password']) + server.ssh_port = 22 + server.private_key = '' + server.public_key = '' + server.save(update_fields=['root_password', 'ssh_port', 'private_key', 'public_key']) passwords[server.address] = password else: errors[server.address] = result.error diff --git a/src/App/tasks.py b/src/App/tasks.py index 2ad1306..4a7d9f9 100644 --- a/src/App/tasks.py +++ b/src/App/tasks.py @@ -12,6 +12,7 @@ import datetime import logging import secrets import threading +import time import dramatiq from django.db import connection @@ -53,6 +54,7 @@ def start(actor, point: Point, status: str, chat, *args): token = secrets.token_hex(16) Point.objects.start_task(point, status, token, chat) actor.send(point.pk, token, *args) + logger.info('Point %s: %s queued, task %s', point, actor.actor_name, token[:8]) supervise.send_with_options(args=(point.pk, token), delay=SUPERVISE_DELAY_MS) @@ -67,15 +69,20 @@ def supervise(point_id: int, token: str): logger.error('Point %s task is stuck: no heartbeat since %s', point, point.task_heartbeat) _abort(point_id, token) return + logger.debug('Point %s task %s alive, last heartbeat %s', point, token[:8], point.task_heartbeat) supervise.send_with_options(args=(point_id, token), delay=SUPERVISE_DELAY_MS) -def _run(point_id: int, token: str, work): - """Выполняет work(point) как задачу token: с отметками для супервизора и снятием, если она прервется.""" +def _run(what: str, point_id: int, token: str, work): + """Выполняет work(point) как задачу token: с отметками для супервизора и снятием, если она прервется. + what - название задачи для лога.""" if not Point.objects.claim_task(point_id, token): # точка удалена, задачу уже снял супервизор или это повторная доставка после падения воркера + logger.warning('Point %s: %s task %s is not claimable, dropped', point_id, what, token[:8]) _abort(point_id, token) return + logger.info('Point %s: %s started, task %s', point_id, what, token[:8]) + started = time.monotonic() stop = threading.Event() heartbeat = threading.Thread(target=_heartbeat, args=(point_id, token, stop), daemon=True) heartbeat.start() @@ -90,6 +97,11 @@ def _run(point_id: int, token: str, work): stop.set() heartbeat.join() Point.objects.finish_task(point_id, token) + point = Point.objects.filter(pk=point_id).first() + logger.info( + 'Point %s: %s finished in %.0fs, %s', point_id, what, time.monotonic() - started, + f'status {point.status}' + (f': {point.error}' if point.error else '') if point else 'point deleted', + ) def _heartbeat(point_id: int, token: str, stop: threading.Event): @@ -138,29 +150,29 @@ def _release(point: Point): # повтор не нужен: пользователь сам решает, переустанавливать ли после ошибки @dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) def install_point(point_id: int, token: str): - _run(point_id, token, lambda point: _install(point, services.install_point)) + _run('install', point_id, token, lambda point: _install(point, services.install_point)) @dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) def reinstall_point(point_id: int, token: str, address: str | None = None, password: str | None = None): - _run(point_id, token, lambda point: _install(point, services.reinstall, address, password)) + _run('reinstall', point_id, token, lambda point: _install(point, services.reinstall, address, password)) @dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) def release_point(point_id: int, token: str): - _run(point_id, token, _release) + _run('release', point_id, token, _release) @dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) def set_domain(point_id: int, token: str, domain: str): - _run(point_id, token, lambda point: _change(point, services.set_domain, domain, texts.domain_changed)) + _run('set_domain', point_id, token, lambda point: _change(point, services.set_domain, domain, texts.domain_changed)) @dramatiq.actor(max_retries=0, time_limit=TIME_LIMIT_MS) def set_outbound_proxy(point_id: int, token: str, encrypted_proxy: str): """encrypted_proxy - fields.encrypt(proxy): в прокси логин и пароль, а сообщения брокера лежат на диске.""" proxy = fields.decrypt(encrypted_proxy) if encrypted_proxy else '' - _run(point_id, token, lambda point: _change(point, services.set_outbound_proxy, proxy, texts.proxy_changed)) + _run('set_outbound_proxy', point_id, token, lambda point: _change(point, services.set_outbound_proxy, proxy, texts.proxy_changed)) # ежедневная проверка: своя цепочка у каждой точки, живет, пока жива точка @@ -170,6 +182,7 @@ def schedule_checks(point: Point): token = secrets.token_hex(16) Point.objects.filter(pk=point.pk).update(check_token=token) check_point.send(point.pk, token) + logger.info('Point %s: daily checks scheduled, chain %s', point, token[:8]) @dramatiq.actor(max_retries=3) @@ -183,10 +196,14 @@ def check_point(point_id: int, token: str): due = point.checked_at is None or timezone.now() - point.checked_at >= CHECK_INTERVAL if point.status == Point.Status.READY and due: _check(point) + else: + logger.debug('Point %s: check skipped, status %s, checked at %s', point, point.status, point.checked_at) def _check(point: Point): + logger.info('Point %s: daily check started', point) problem = services.check_point(point) + logger.info('Point %s: daily check %s', point, f'found problem: {problem}' if problem else 'passed') previous = point.health_error Point.objects.filter(pk=point.pk).update(checked_at=timezone.now(), health_error=problem) if problem == previous: diff --git a/src/App/tests_services.py b/src/App/tests_services.py index db6a1ee..05e83dc 100644 --- a/src/App/tests_services.py +++ b/src/App/tests_services.py @@ -12,15 +12,17 @@ from .models import Credentials, Point, Server class FakeRunner: - """Подменяет запуск плейбуков serverus. fail - плейбук, который завершится ошибкой.""" + """Подменяет запуск плейбуков serverus. fail - плейбук, который завершится ошибкой, + fail_address - только на этом сервере.""" - def __init__(self, fail: str | None = None): + def __init__(self, fail: str | None = None, fail_address: str | None = None): self.calls = [] self.fail = fail + self.fail_address = fail_address def __call__(self, playbook, host, extravars): self.calls.append((playbook, host, extravars)) - if playbook == self.fail: + if playbook == self.fail and self.fail_address in (None, host.address): return serverus.Result(ok=False, status='failed', error='boom', known_hosts=f'kh {host.address}') data = { 'facts': {'uname': 'Linux test'}, @@ -431,3 +433,20 @@ class ReleaseTests(TestCase): self.assertEqual(set(errors), {'10.0.0.1', '10.0.0.2'}) self.point.refresh_from_db() self.assertEqual(self.point.status, Point.Status.ERROR) + + def test_retry_skips_released_server(self): + # панель вернулась, маскировка нет: повтор не должен идти на панель со старым ключом и портом + runner = FakeRunner(fail='release', fail_address='10.0.0.2') + with patch_runner(runner): + first, errors = services.release_point(self.point) + self.assertEqual(set(errors), {'10.0.0.2'}) + panel = self.point.servers.get(address='10.0.0.1') + self.assertEqual((panel.ssh_port, panel.private_key), (22, '')) + + runner = FakeRunner() + with patch_runner(runner): + passwords, errors = services.release_point(self.point) + self.assertEqual(errors, {}) + self.assertEqual(runner.playbooks('10.0.0.1'), []) + self.assertEqual(passwords['10.0.0.1'], first['10.0.0.1']) + self.assertFalse(Point.objects.exists()) diff --git a/src/MyPointVPN/settings/README.md b/src/MyPointVPN/settings/README.md index 7ecffcf..5446266 100644 --- a/src/MyPointVPN/settings/README.md +++ b/src/MyPointVPN/settings/README.md @@ -12,4 +12,4 @@ from .project import * ``` - `test.py` - настройки для тестов: фоновые задачи идут в stub брокер, а не в rabbitmq. `make dj_test` использует их. -- `docker.py` - настройки для docker: все из окружения, `DEBUG=False`, `SECRET_KEY` обязателен, статика через whitenoise, логи в консоль. +- `docker.py` - настройки для docker: все из окружения, `DEBUG=False`, `SECRET_KEY` обязателен, статика через whitenoise. diff --git a/src/MyPointVPN/settings/docker.py b/src/MyPointVPN/settings/docker.py index 80717e5..7078474 100644 --- a/src/MyPointVPN/settings/docker.py +++ b/src/MyPointVPN/settings/docker.py @@ -33,15 +33,3 @@ STORAGES = { 'default': {'BACKEND': 'django.core.files.storage.FileSystemStorage'}, 'staticfiles': {'BACKEND': 'whitenoise.storage.CompressedManifestStaticFilesStorage'}, } - -LOGGING = { - 'version': 1, - 'disable_existing_loggers': False, - 'formatters': { - 'plain': {'format': '{asctime} {levelname} {name}: {message}', 'style': '{'}, - }, - 'handlers': { - 'console': {'class': 'logging.StreamHandler', 'formatter': 'plain'}, - }, - 'root': {'handlers': ['console'], 'level': os.environ.get('LOG_LEVEL', 'INFO')}, -} diff --git a/src/MyPointVPN/settings/project.py b/src/MyPointVPN/settings/project.py index 0295432..f19cf77 100644 --- a/src/MyPointVPN/settings/project.py +++ b/src/MyPointVPN/settings/project.py @@ -72,6 +72,23 @@ SERVERUS_CA_DIR = os.environ.get('SERVERUS_CA_DIR', str(BASE_DIR.parent / 'extra # Бэкапы базы: manage.py backup и restore BACKUP_DIR = os.environ.get('BACKUP_DIR') or str(BASE_DIR.parent / 'extra/backups') +# Логи в консоль: сайт, воркер и docker. pid и поток - чтобы различать процессы и задачи воркера dramatiq +LOGGING = { + 'version': 1, + 'disable_existing_loggers': False, + 'formatters': { + 'plain': {'format': '{asctime} {levelname} [{process}:{threadName}] {name}: {message}', 'style': '{'}, + }, + 'handlers': { + 'console': {'class': 'logging.StreamHandler', 'formatter': 'plain'}, + }, + 'root': {'handlers': ['console'], 'level': os.environ.get('LOG_LEVEL', 'INFO')}, + 'loggers': { + # соединения с брокером на каждую отправку задачи + 'pika': {'level': 'WARNING'}, + }, +} + # Фоновые задачи: dramatiq поверх rabbitmq. Воркер: make dj_worker RABBITMQ_URL = 'amqp://{user}:{password}@{address}/{vhost}'.format( user=quote(os.environ.get('RABBITMQ_USER', 'guest'), safe=''), diff --git a/src/serverus/playbooks/release.yml b/src/serverus/playbooks/release.yml index 956926b..91257c1 100644 --- a/src/serverus/playbooks/release.yml +++ b/src/serverus/playbooks/release.yml @@ -53,10 +53,13 @@ changed_when: false failed_when: false - # systemd-resolved остается: без него /etc/resolv.conf указывает в никуда + # systemd-resolved остается: без него /etc/resolv.conf указывает в никуда. + # DNS настраивается только на панели, на сервере маскировки resolved может не быть: rc 5 - юнита нет - name: Restart systemd-resolved ansible.builtin.command: systemctl try-restart systemd-resolved + register: resolved_restart changed_when: false + failed_when: resolved_restart.rc not in [0, 5] - name: Remove sshd drop-in ansible.builtin.file: diff --git a/src/serverus/runner.py b/src/serverus/runner.py index 9f776e2..fdfc81f 100644 --- a/src/serverus/runner.py +++ b/src/serverus/runner.py @@ -1,6 +1,8 @@ +import logging import os import sys import tempfile +import time from dataclasses import dataclass, field from pathlib import Path from typing import Protocol @@ -9,6 +11,8 @@ import ansible_runner from .host import Host +logger = logging.getLogger(__name__) + PLAYBOOKS_DIR = Path(__file__).resolve().parent / 'playbooks' TARGET = 'target' # ansible-playbook ставится рядом с интерпретатором, но venv может быть не активирован @@ -35,6 +39,8 @@ def run_playbook(playbook: str, host: Host, extravars: dict) -> Result: Все файлы запуска (ключ, extravars, артефакты) живут во временной папке и удаляются после выполнения. """ + logger.info('%s: playbook %s started', host.address, playbook) + started = time.monotonic() with tempfile.TemporaryDirectory(prefix='serverus-') as private_dir: result = ansible_runner.run( private_data_dir=private_dir, @@ -51,21 +57,36 @@ def run_playbook(playbook: str, host: Host, extravars: dict) -> Result: 'ANSIBLE_HOST_KEY_CHECKING': 'True', }, quiet=True, + event_handler=lambda event: _log_event(host, playbook, event), ) events = list(result.events) with result.stdout as stdout: output = stdout.read() known_hosts = _known_hosts_file(private_dir).read_text() ok = result.status == 'successful' + error = '' if ok else _error(events) or output.strip()[-1000:] + elapsed = time.monotonic() - started + if ok: + logger.info('%s: playbook %s finished in %.0fs', host.address, playbook, elapsed) + else: + logger.warning('%s: playbook %s %s in %.0fs: %s', host.address, playbook, result.status, elapsed, error) return Result( ok=ok, status=result.status, data=_stats_data(events), - error='' if ok else _error(events) or output.strip()[-1000:], + error=error, known_hosts=known_hosts, ) +def _log_event(host: Host, playbook: str, event: dict) -> bool: + """Шаги плейбука в лог: установка идет до часа, так видно, на чем она сейчас.""" + if event.get('event') == 'playbook_on_task_start': + logger.info('%s: %s: %s', host.address, playbook, event.get('event_data', {}).get('task', '')) + # True - событие остается в result.events + return True + + def _known_hosts_file(private_dir: str) -> Path: return Path(private_dir) / 'known_hosts' @@ -112,5 +133,9 @@ def _error(events: list[dict]) -> str: if event.get('event') in ('runner_on_failed', 'runner_on_unreachable'): event_data = event.get('event_data', {}) res = event_data.get('res', {}) - return f"{event_data.get('task', '')}: {res.get('msg') or res.get('stderr') or 'failed'}" + if event['event'] == 'runner_on_unreachable': + # у задач с no_log ответ скрыт целиком, включая причину недоступности + return f"{event_data.get('task', '')}: сервер недоступен: {res.get('msg') or 'ssh'}" + # msg у command - только "non-zero return code", причина в stderr + return f"{event_data.get('task', '')}: {res.get('stderr') or res.get('msg') or 'failed'}" return ''