92 lines
3.1 KiB
Python
92 lines
3.1 KiB
Python
"""
|
||
Database connection — psycopg2 connection pool.
|
||
|
||
Архитектурное решение (Opus):
|
||
- ThreadedConnectionPool(minconn=1, maxconn=10) — оптимально для ThreadingMixIn HTTP-сервера.
|
||
Каждый HTTP-запрос в отдельном потоке получает своё соединение из пула.
|
||
- RealDictCursor для query() — возвращает dict с lowercase ключами (консистентно с Python API).
|
||
- execute()/execute_returning() — для INSERT/UPDATE/DELETE с автокоммитом и откатом при ошибке.
|
||
- Пул создаётся лениво (при первом запросе) через get_pool().
|
||
"""
|
||
import os
|
||
import psycopg2
|
||
import psycopg2.pool
|
||
import psycopg2.extras
|
||
|
||
# Глобальный пул соединений (singleton)
|
||
_pool = None
|
||
|
||
# Параметры подключения из переменных окружения (systemd Environment)
|
||
DB_CONFIG = {
|
||
"host": os.getenv("DB_HOST", "127.0.0.1"),
|
||
"port": int(os.getenv("DB_PORT", "5432")),
|
||
"dbname": os.getenv("DB_NAME", "baza"),
|
||
"user": os.getenv("DB_USER", "super"),
|
||
"password": os.getenv("DB_PASS", ""),
|
||
}
|
||
|
||
|
||
def get_pool():
|
||
"""Возвращает глобальный пул соединений. Создаёт при первом вызове."""
|
||
global _pool
|
||
if _pool is None:
|
||
_pool = psycopg2.pool.ThreadedConnectionPool(
|
||
minconn=1, maxconn=10, **DB_CONFIG
|
||
)
|
||
return _pool
|
||
|
||
|
||
def query(sql, params=None):
|
||
"""
|
||
SELECT → list[dict] с lowercase ключами.
|
||
Используется всеми db/*.py модулями для чтения данных.
|
||
"""
|
||
pool = get_pool()
|
||
conn = pool.getconn()
|
||
try:
|
||
with conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur:
|
||
cur.execute(sql, params)
|
||
rows = cur.fetchall()
|
||
return [{k.lower(): v for k, v in r.items()} for r in rows]
|
||
finally:
|
||
pool.putconn(conn)
|
||
|
||
|
||
def execute(sql, params=None):
|
||
"""
|
||
INSERT/UPDATE/DELETE → количество затронутых строк.
|
||
Автокоммит. При ошибке — rollback и проброс исключения.
|
||
"""
|
||
pool = get_pool()
|
||
conn = pool.getconn()
|
||
try:
|
||
with conn.cursor() as cur:
|
||
cur.execute(sql, params)
|
||
conn.commit()
|
||
return cur.rowcount
|
||
except Exception:
|
||
conn.rollback()
|
||
raise
|
||
finally:
|
||
pool.putconn(conn)
|
||
|
||
|
||
def execute_returning(sql, params=None):
|
||
"""
|
||
INSERT/UPDATE/DELETE с RETURNING → dict первой строки.
|
||
Используется для insert с автогенерацией UUID (gen_random_uuid()).
|
||
"""
|
||
pool = get_pool()
|
||
conn = pool.getconn()
|
||
try:
|
||
with conn.cursor(cursor_factory=psycopg2.extras.RealDictCursor) as cur:
|
||
cur.execute(sql, params)
|
||
conn.commit()
|
||
row = cur.fetchone()
|
||
return {k.lower(): v for k, v in row.items()} if row else None
|
||
except Exception:
|
||
conn.rollback()
|
||
raise
|
||
finally:
|
||
pool.putconn(conn)
|