fix: не мог вернуть сервер маскировки

This commit is contained in:
protokey committed 2026-09-26 19:39:44 +04:00
1 parent f1414e5ace
commit 65c19ea1cd
9 files changed
+106 -27

No files matched your search

+2
View File
@@ -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
+9 -1
View File
@@ -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
+24 -7
View File
@@ -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:
+22 -3
View File
@@ -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())
+1 -1
View File
@@ -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.
-12
View File
@@ -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')},
}
+17
View File
@@ -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=''),
+4 -1
View File
@@ -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:
+27 -2
View File
@@ -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 ''