186 lines
8.5 KiB
Python
186 lines
8.5 KiB
Python
"""
|
||
Трекер созданных инстансов — краткосрочный 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())
|