- execute(sql, params) → (rowcount, error) — для INSERT/UPDATE/DELETE - query_one(sql, params) → (dict, error) — первая строка как словарь - авто-commit и rollback при ошибке
104 lines
3.3 KiB
Python
104 lines
3.3 KiB
Python
"""
|
|
db.py — Транспортный слой к базе данных.
|
|
|
|
Только connect() и query(). Никакой бизнес-логики, никакого DDL.
|
|
DDL → schema.py. Бизнес-запросы → upload.py, extractor.py, api.py.
|
|
|
|
Функции:
|
|
_pg_connect(dbname) — сырое подключение к ЛЮБОЙ базе (используется /test createdb)
|
|
connect() — подключение к ЦЕЛЕВОЙ базе (из DB_NAME)
|
|
query(sql, params) — выполнить SELECT → (result, error)
|
|
execute(sql, params) — выполнить INSERT/UPDATE/DELETE → (rowcount, error)
|
|
query_one(sql, params) — как query, но возвращает одну строку или None
|
|
"""
|
|
|
|
import os
|
|
import psycopg2
|
|
|
|
|
|
def _pg_connect(dbname):
|
|
"""
|
|
Сырое подключение к указанной базе данных.
|
|
Используется:
|
|
- connect() для целевой БД
|
|
- test_routes.py для /test createdb (подключение к 'postgres')
|
|
"""
|
|
return psycopg2.connect(
|
|
host=os.getenv("DB_HOST"),
|
|
port=os.getenv("DB_PORT", 5432),
|
|
dbname=dbname,
|
|
user=os.getenv("DB_USER"),
|
|
password=os.getenv("DB_PASS"),
|
|
sslmode=os.getenv("DB_SSLMODE", "disable"),
|
|
)
|
|
|
|
|
|
def connect():
|
|
"""
|
|
Подключение к ЦЕЛЕВОЙ базе данных (DB_NAME из переменных окружения).
|
|
Возвращает (connection, None) при успехе или (None, error) при ошибке.
|
|
"""
|
|
try:
|
|
return _pg_connect(os.getenv("DB_NAME")), None
|
|
except Exception as e:
|
|
return None, str(e)
|
|
|
|
|
|
def query(sql_text, params=None):
|
|
"""
|
|
Выполнить SQL-запрос к целевой БД.
|
|
Возвращает (result, None) или (None, error).
|
|
|
|
result = {"columns": [...], "rows": [[...], ...]}
|
|
"""
|
|
conn, err = connect()
|
|
if err:
|
|
return None, f"connect: {err}"
|
|
try:
|
|
cur = conn.cursor()
|
|
cur.execute(sql_text, params)
|
|
rows = cur.fetchall()
|
|
cols = [desc[0] for desc in cur.description] if cur.description else []
|
|
cur.close()
|
|
conn.close()
|
|
return {"columns": cols, "rows": [list(r) for r in rows]}, None
|
|
except Exception as e:
|
|
conn.close()
|
|
return None, str(e)
|
|
|
|
|
|
def execute(sql_text, params=None):
|
|
"""
|
|
Выполнить INSERT/UPDATE/DELETE к целевой БД.
|
|
Возвращает (rowcount, None) или (None, error).
|
|
"""
|
|
conn, err = connect()
|
|
if err:
|
|
return None, f"connect: {err}"
|
|
try:
|
|
cur = conn.cursor()
|
|
cur.execute(sql_text, params)
|
|
conn.commit()
|
|
rowcount = cur.rowcount
|
|
cur.close()
|
|
conn.close()
|
|
return rowcount, None
|
|
except Exception as e:
|
|
conn.rollback()
|
|
conn.close()
|
|
return None, str(e)
|
|
|
|
|
|
def query_one(sql_text, params=None):
|
|
"""
|
|
Выполнить SELECT и вернуть ПЕРВУЮ строку или None.
|
|
Возвращает (row_dict, None) или (None, error).
|
|
"""
|
|
result, err = query(sql_text, params)
|
|
if err:
|
|
return None, err
|
|
rows = result["rows"]
|
|
if not rows:
|
|
return None, None # нет строк — не ошибка
|
|
return dict(zip(result["columns"], rows[0])), None
|