71 lines
1.8 KiB
Python
71 lines
1.8 KiB
Python
"""Database connection — psycopg2 connection pool."""
|
|
import os
|
|
import psycopg2
|
|
import psycopg2.pool
|
|
import psycopg2.extras
|
|
|
|
_pool = None
|
|
|
|
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 of dicts (lowercase keys)."""
|
|
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 → rowcount."""
|
|
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 with RETURNING → first row dict."""
|
|
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)
|