Files

186 lines
8.5 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
Трекер созданных инстансов — краткосрочный JSON-файловый кеш.
ПРОБЛЕМА: когда мы создаём инстанс через POST /instances, он появляется
в Nubes НЕ МГНОВЕННО. GET /instances может задержать новый инстанс на 5-30 секунд.
В это время пользователь не видит свой инстанс в списке → думает что create не сработал.
РЕШЕНИЕ: трекер — это JSON-файл в /tmp/, куда мы пишем UID инстанса СРАЗУ после
успешного create. UI сначала ищет инстансы в облаке (cloud-first), а то чего
нет в облаке — добирает из трекера (tracker-fallback).
Файлы изолированы по пользователю и стенду:
/tmp/instances-{clientId}-{stand}.json
Блокировка fcntl.flock (эксклюзивная, неблокирующая с retry) —
для безопасной конкурентной работы нескольких gunicorn-воркеров.
ЖИЗНЕННЫЙ ЦИКЛ:
- add() — executor вызывает сразу после получения instanceUid
- remove() — после успешного delete (в _finish_op и CMDB-delete)
- При редеплое /tmp/ теряется — это НЕ критично (cloud-first всё покажет)
Функции:
add(client_id, stand, uid, svc_id, name) — добавить инстанс
remove(client_id, stand, uid) — удалить (после успешного delete)
list_all(client_id, stand) — все записи списком {svcId, displayName, instanceUid}
"""
import fcntl
import json
import os
import time
def _path(client_id, stand):
"""Путь к файлу трекера: /tmp/instances-{clientId}-{stand}.json.
Изоляция по clientId + stand гарантирует что пользователи не видят чужие инстансы."""
return f"/tmp/instances-{client_id}-{stand}.json"
def _acquire_lock(fd):
"""Эксклюзивная блокировка файла с таймаутом 2 секунды.
Использует LOCK_NB (неблокирующий) + retry с шагом 50ms.
Почему не LOCK_EX с бесконечным ожиданием:
- При подвисании воркера все остальные повиснут навсегда.
- Лучше не записать в трекер, чем уронить запрос.
Returns:
True — блокировка взята.
False — не смогли за 2 секунды (другой воркер держит слишком долго)."""
deadline = time.time() + 2
while True:
try:
fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
return True # блокировка взята
except BlockingIOError:
if time.time() >= deadline:
return False # таймаут — сдаёмся
time.sleep(0.05)
def _locked_read(path):
"""Читает JSON-файл под эксклюзивной блокировкой.
Если файла нет — создаёт (O_CREAT) и возвращает {}.
Если JSON битый — возвращает {} (переживёт перезапись при следующем add).
Максимальный размер чтения: 65536 байт (64KB) — больше трекеру не нужно.
Returns:
dict — содержимое файла или {}."""
try:
# O_RDWR — чтение+запись, O_CREAT — создать если нет
fd = os.open(path, os.O_RDWR | os.O_CREAT, 0o644)
except OSError:
return {} # нет прав (например /tmp/ только read) — молча
try:
if not _acquire_lock(fd):
return {} # не смогли заблокировать → не читаем (гонка данных)
try:
data = os.read(fd, 65536)
if data:
return json.loads(data.decode("utf-8"))
except (json.JSONDecodeError, UnicodeDecodeError):
pass # битый JSON — при следующей записи перезапишется
return {}
finally:
fcntl.flock(fd, fcntl.LOCK_UN) # ВСЕГДА снимаем блокировку
os.close(fd)
def _locked_write(path, data):
"""Пишет JSON-файл под эксклюзивной блокировкой.
Полностью перезаписывает файл: lseek(0) + ftruncate + write.
Это атомарно под локом — другие воркеры увидят либо старую, либо новую версию.
Если не смогли взять лок — НЕ пишем.
Данные не потеряются — инстанс создан в Nubes, cloud-first его подхватит."""
try:
fd = os.open(path, os.O_RDWR | os.O_CREAT, 0o644)
except OSError:
return
try:
if not _acquire_lock(fd):
return # не смогли заблокировать — пропускаем запись
# Три шага атомарной перезаписи:
os.lseek(fd, 0, 0) # 1. В начало файла
os.ftruncate(fd, 0) # 2. Обрезать старый контент
os.write(fd, json.dumps(data, indent=2).encode("utf-8")) # 3. Записать новый
finally:
fcntl.flock(fd, fcntl.LOCK_UN) # ВСЕГДА снимаем блокировку
os.close(fd)
def _atomic_update(path, mutator):
"""Атомарно: открыть файл → заблокировать → прочитать → мутировать → записать.
Вся операция под ОДНИМ lock, в отличие от старого подхода
_locked_read() + _locked_write() где между ними lock снимался.
Это предотвращает lost-update при конкурентных add/remove.
Args:
path: str — путь к файлу
mutator: callable(data) — функция, мутирующая dict (возвращать не нужно)
Returns:
dict — результат после мутации (или {} если не смогли)"""
try:
fd = os.open(path, os.O_RDWR | os.O_CREAT, 0o644)
except OSError:
return {}
try:
if not _acquire_lock(fd):
return {}
# Читаем
data = {}
try:
raw = os.read(fd, 65536)
if raw:
data = json.loads(raw.decode("utf-8"))
except (json.JSONDecodeError, UnicodeDecodeError):
pass
# Мутируем
mutator(data)
# Пишем
os.lseek(fd, 0, 0)
os.ftruncate(fd, 0)
os.write(fd, json.dumps(data, indent=2).encode("utf-8"))
return data
finally:
fcntl.flock(fd, fcntl.LOCK_UN)
os.close(fd)
def add(client_id, stand, instance_uid, svc_id, display_name):
"""Добавить инстанс в трекер (атомарно, под одним lock).
Читает → добавляет/обновляет запись → пишет — всё под эксклюзивной блокировкой.
Если инстанс уже есть — перезаписывает (идемпотентность)."""
def _add(data):
data[instance_uid] = {
"svcId": svc_id,
"displayName": display_name,
"instanceUid": instance_uid,
}
_atomic_update(_path(client_id, stand), _add)
def remove(client_id, stand, instance_uid):
"""Удалить инстанс из трекера (атомарно, под одним lock).
pop с default=None — не падает если инстанса уже нет."""
def _rem(data):
data.pop(instance_uid, None)
_atomic_update(_path(client_id, stand), _rem)
def list_all(client_id, stand):
"""Все записи трекера — список dict'ов {svcId, displayName, instanceUid}.
Используется в main.py → api_operations для tracker-fallback.
Конвертирует dict (indexed by UUID) в list."""
return list(_locked_read(_path(client_id, stand)).values())