63 lines
1.9 KiB
Python
63 lines
1.9 KiB
Python
import os
|
|
import psycopg2
|
|
from psycopg2 import sql
|
|
|
|
|
|
def _pg_connect(dbname):
|
|
"""Сырое подключение к указанной базе."""
|
|
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():
|
|
"""Подключение к целевой БД. Возвращает (connection, None) или (None, error)."""
|
|
try:
|
|
return _pg_connect(os.getenv("DB_NAME")), None
|
|
except Exception as e:
|
|
return None, str(e)
|
|
|
|
|
|
def ensure_db():
|
|
"""Создать базу, если не существует."""
|
|
db_name = os.getenv("DB_NAME", "contracts")
|
|
try:
|
|
conn = _pg_connect("postgres")
|
|
conn.autocommit = True
|
|
cur = conn.cursor()
|
|
cur.execute("SELECT 1 FROM pg_database WHERE datname = %s", (db_name,))
|
|
if cur.fetchone():
|
|
print(f"DB '{db_name}' exists")
|
|
else:
|
|
cur.execute(sql.SQL("CREATE DATABASE {}").format(sql.Identifier(db_name)))
|
|
print(f"DB '{db_name}' created")
|
|
cur.close()
|
|
conn.close()
|
|
return True
|
|
except Exception as e:
|
|
print(f"ensure_db error: {e}")
|
|
return False
|
|
|
|
|
|
def query(sql_text, params=None):
|
|
"""Выполнить запрос, вернуть (result, error)."""
|
|
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)
|