1
0
Fork 0
forked from zovos/bot_tg
bot_tg/scripts/migrate_sqlite_to_pg.py
q add3553ce0 feat: переход с SQLite на PostgreSQL + скрипт миграции
- Все 4 SQLite базы объединены в одну PostgreSQL
- Новый db.py с ThreadedConnectionPool (psycopg2)
- docker-compose: сервис postgres с healthcheck, данные в ./db/postgres
- scripts/migrate_sqlite_to_pg.py — миграция существующих данных
- Убраны мёртвый код и дублирующий import в main.py и talk_handler.py
- .env удалён из git-индекса

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-01 23:25:22 +03:00

369 lines
13 KiB
Python

#!/usr/bin/env python3
"""
Migrate data from SQLite databases to PostgreSQL.
Usage (from repo root):
docker compose run --rm bot python3 scripts/migrate_sqlite_to_pg.py
Or locally (PostgreSQL must be reachable):
DATABASE_URL=postgresql://bot:botpass@localhost:5432/botdb python3 scripts/migrate_sqlite_to_pg.py
SQLite sources expected in ./db/:
economy.sqlite3 → user_balances, deposits, loans, transactions, taxes,
economy_settings, daily_activity, central_bank, penis_stats
bets.sqlite3 → bets
chat_history.sqlite3→ chat_history, chat_state, talk_history, talk_state
penis_stats.sqlite3 → (legacy) length column → user_balances.balance + penis_stats
"""
import os
import sqlite3
import sys
from pathlib import Path
import psycopg2
import psycopg2.extras
DB_DIR = Path(os.getenv("SQLITE_DIR", "/db"))
DATABASE_URL = os.getenv("DATABASE_URL", "postgresql://bot:botpass@postgres:5432/botdb")
ECONOMY_DB = DB_DIR / "economy.sqlite3"
BETS_DB = DB_DIR / "bets.sqlite3"
CHAT_DB = DB_DIR / "chat_history.sqlite3"
PENIS_LEGACY_DB = DB_DIR / "penis_stats.sqlite3"
def open_sqlite(path: Path) -> sqlite3.Connection | None:
if not path.exists():
print(f" [skip] {path} not found")
return None
conn = sqlite3.connect(path)
conn.row_factory = sqlite3.Row
return conn
def migrate_economy(pg: psycopg2.extensions.connection) -> None:
src = open_sqlite(ECONOMY_DB)
if src is None:
return
cur = pg.cursor()
# --- user_balances ---
rows = src.execute("SELECT * FROM user_balances").fetchall()
print(f" user_balances: {len(rows)} rows")
for r in rows:
cur.execute("""
INSERT INTO user_balances (user_id, balance, daily_income, last_daily_reset, created_at)
VALUES (%s, %s, %s, %s, %s)
ON CONFLICT (user_id) DO UPDATE SET
balance = EXCLUDED.balance,
daily_income = EXCLUDED.daily_income,
last_daily_reset = EXCLUDED.last_daily_reset
""", (
r["user_id"],
r["balance"],
r["daily_income"],
r["last_daily_reset"] or None,
r["created_at"] or None,
))
# --- deposits ---
rows = src.execute("SELECT * FROM deposits").fetchall()
print(f" deposits: {len(rows)} rows")
for r in rows:
cur.execute("""
INSERT INTO deposits (id, user_id, amount, interest_rate, created_at, matures_at, is_active)
VALUES (%s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (id) DO NOTHING
""", (
r["id"], r["user_id"], r["amount"], r["interest_rate"],
r["created_at"] or None, r["matures_at"] or None,
bool(r["is_active"]),
))
_reset_sequence(cur, "deposits", "id")
# --- loans ---
rows = src.execute("SELECT * FROM loans").fetchall()
print(f" loans: {len(rows)} rows")
for r in rows:
cur.execute("""
INSERT INTO loans (id, user_id, amount, interest_rate, created_at, due_at, is_repaid)
VALUES (%s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (id) DO NOTHING
""", (
r["id"], r["user_id"], r["amount"], r["interest_rate"],
r["created_at"] or None, r["due_at"] or None,
bool(r["is_repaid"]),
))
_reset_sequence(cur, "loans", "id")
# --- transactions ---
rows = src.execute("SELECT * FROM transactions").fetchall()
print(f" transactions: {len(rows)} rows")
for r in rows:
cur.execute("""
INSERT INTO transactions (id, from_user_id, to_user_id, amount, transaction_type, description, created_at)
VALUES (%s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (id) DO NOTHING
""", (
r["id"], r["from_user_id"], r["to_user_id"], r["amount"],
r["transaction_type"], r["description"], r["created_at"] or None,
))
_reset_sequence(cur, "transactions", "id")
# --- taxes ---
rows = src.execute("SELECT * FROM taxes").fetchall()
print(f" taxes: {len(rows)} rows")
for r in rows:
cur.execute("""
INSERT INTO taxes (id, user_id, amount, tax_date, created_at)
VALUES (%s, %s, %s, %s, %s)
ON CONFLICT (id) DO NOTHING
""", (
r["id"], r["user_id"], r["amount"],
r["tax_date"] or None, r["created_at"] or None,
))
_reset_sequence(cur, "taxes", "id")
# --- economy_settings ---
rows = src.execute("SELECT * FROM economy_settings").fetchall()
print(f" economy_settings: {len(rows)} rows")
for r in rows:
cur.execute("""
INSERT INTO economy_settings (key, value, updated_at)
VALUES (%s, %s, %s)
ON CONFLICT (key) DO UPDATE SET value = EXCLUDED.value
""", (r["key"], r["value"], r["updated_at"] if "updated_at" in r.keys() else None))
# --- daily_activity ---
rows = src.execute("SELECT * FROM daily_activity").fetchall()
print(f" daily_activity: {len(rows)} rows")
for r in rows:
cur.execute("""
INSERT INTO daily_activity (id, user_id, activity_date, message_count, last_activity_time, created_at)
VALUES (%s, %s, %s, %s, %s, %s)
ON CONFLICT (user_id, activity_date) DO UPDATE SET
message_count = EXCLUDED.message_count,
last_activity_time = EXCLUDED.last_activity_time
""", (
r["id"], r["user_id"],
r["activity_date"] or None, r["message_count"],
r["last_activity_time"] or None, r["created_at"] or None,
))
_reset_sequence(cur, "daily_activity", "id")
# --- central_bank ---
rows = src.execute("SELECT * FROM central_bank").fetchall()
print(f" central_bank: {len(rows)} rows")
for r in rows:
cur.execute("""
INSERT INTO central_bank (id, capital, total_taxes_collected,
total_inflation_adjustments, total_emissions, last_updated)
VALUES (%s, %s, %s, %s, %s, %s)
ON CONFLICT (id) DO UPDATE SET
capital = EXCLUDED.capital,
total_taxes_collected = EXCLUDED.total_taxes_collected,
total_inflation_adjustments = EXCLUDED.total_inflation_adjustments,
total_emissions = EXCLUDED.total_emissions,
last_updated = EXCLUDED.last_updated
""", (
r["id"], r["capital"],
r["total_taxes_collected"], r["total_inflation_adjustments"],
r["total_emissions"], r["last_updated"] or None,
))
# --- penis_stats (new schema: user_id, display_name, last_used_ts) ---
try:
rows = src.execute("SELECT * FROM penis_stats").fetchall()
print(f" penis_stats: {len(rows)} rows")
cols = rows[0].keys() if rows else []
for r in rows:
if "last_used_ts" in cols:
cur.execute("""
INSERT INTO penis_stats (user_id, display_name, last_used_ts)
VALUES (%s, %s, %s)
ON CONFLICT (user_id) DO UPDATE SET
display_name = EXCLUDED.display_name,
last_used_ts = EXCLUDED.last_used_ts
""", (r["user_id"], r["display_name"], r["last_used_ts"]))
else:
# old schema without last_used_ts
cur.execute("""
INSERT INTO penis_stats (user_id, display_name)
VALUES (%s, %s)
ON CONFLICT (user_id) DO UPDATE SET display_name = EXCLUDED.display_name
""", (r["user_id"], r.get("display_name", "")))
except sqlite3.OperationalError as e:
print(f" [warn] penis_stats in economy.sqlite3: {e}")
src.close()
pg.commit()
print(" economy.sqlite3 done")
def migrate_legacy_penis(pg: psycopg2.extensions.connection) -> None:
"""Old penis_stats.sqlite3 had a `length` column — map it to user_balances.balance."""
src = open_sqlite(PENIS_LEGACY_DB)
if src is None:
return
cur = pg.cursor()
try:
rows = src.execute("SELECT user_id, length, display_name FROM penis_stats").fetchall()
except sqlite3.OperationalError:
rows = src.execute("SELECT * FROM penis_stats").fetchall()
print(f" penis_stats.sqlite3 (legacy): {len(rows)} rows")
for r in rows:
uid = r["user_id"]
length = r["length"] if "length" in r.keys() else 0.0
name = r["display_name"] if "display_name" in r.keys() else ""
last_ts = r["last_used_ts"] if "last_used_ts" in r.keys() else None
# upsert into user_balances; don't overwrite if already migrated from economy.sqlite3
cur.execute("""
INSERT INTO user_balances (user_id, balance)
VALUES (%s, %s)
ON CONFLICT (user_id) DO NOTHING
""", (uid, float(length)))
cur.execute("""
INSERT INTO penis_stats (user_id, display_name, last_used_ts)
VALUES (%s, %s, %s)
ON CONFLICT (user_id) DO UPDATE SET
display_name = EXCLUDED.display_name,
last_used_ts = COALESCE(EXCLUDED.last_used_ts, penis_stats.last_used_ts)
""", (uid, name, last_ts))
src.close()
pg.commit()
print(" penis_stats.sqlite3 done")
def migrate_bets(pg: psycopg2.extensions.connection) -> None:
src = open_sqlite(BETS_DB)
if src is None:
return
cur = pg.cursor()
rows = src.execute("SELECT * FROM bets").fetchall()
print(f" bets: {len(rows)} rows")
cols = rows[0].keys() if rows else []
for r in rows:
cur.execute("""
INSERT INTO bets (id, user_id, match_id, sport, home_team, away_team,
chosen_team, amount, odds, status, payout, created_ts, resolved_ts)
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (id) DO NOTHING
""", (
r["id"], r["user_id"], r["match_id"], r["sport"],
r["home_team"], r["away_team"], r["chosen_team"],
r["amount"], r["odds"], r["status"],
r["payout"] if "payout" in cols else 0,
r["created_ts"] if "created_ts" in cols else None,
r["resolved_ts"] if "resolved_ts" in cols else None,
))
_reset_sequence(cur, "bets", "id")
src.close()
pg.commit()
print(" bets.sqlite3 done")
def migrate_chat(pg: psycopg2.extensions.connection) -> None:
src = open_sqlite(CHAT_DB)
if src is None:
return
cur = pg.cursor()
# --- chat_history ---
rows = src.execute("SELECT * FROM chat_history").fetchall()
print(f" chat_history: {len(rows)} rows")
cols = rows[0].keys() if rows else []
for r in rows:
cur.execute("""
INSERT INTO chat_history (id, chat_id, role, name, text, created_at)
VALUES (%s, %s, %s, %s, %s, %s)
ON CONFLICT (id) DO NOTHING
""", (
r["id"], r["chat_id"], r["role"], r["name"], r["text"],
r["created_at"] if "created_at" in cols else 0,
))
_reset_sequence(cur, "chat_history", "id")
# --- chat_state ---
rows = src.execute("SELECT * FROM chat_state").fetchall()
print(f" chat_state: {len(rows)} rows")
for r in rows:
cur.execute("""
INSERT INTO chat_state (chat_id, key, value)
VALUES (%s, %s, %s)
ON CONFLICT (chat_id, key) DO UPDATE SET value = EXCLUDED.value
""", (r["chat_id"], r["key"], r["value"]))
# --- talk_history ---
rows = src.execute("SELECT * FROM talk_history").fetchall()
print(f" talk_history: {len(rows)} rows")
cols = rows[0].keys() if rows else []
for r in rows:
cur.execute("""
INSERT INTO talk_history (id, chat_id, user_id, role, name, text, created_at)
VALUES (%s, %s, %s, %s, %s, %s, %s)
ON CONFLICT (id) DO NOTHING
""", (
r["id"], r["chat_id"], r["user_id"], r["role"], r["name"], r["text"],
r["created_at"] if "created_at" in cols else 0,
))
_reset_sequence(cur, "talk_history", "id")
# --- talk_state ---
rows = src.execute("SELECT * FROM talk_state").fetchall()
print(f" talk_state: {len(rows)} rows")
for r in rows:
cur.execute("""
INSERT INTO talk_state (chat_id, user_id, key, value)
VALUES (%s, %s, %s, %s)
ON CONFLICT (chat_id, user_id, key) DO UPDATE SET value = EXCLUDED.value
""", (r["chat_id"], r["user_id"], r["key"], r["value"]))
src.close()
pg.commit()
print(" chat_history.sqlite3 done")
def _reset_sequence(cur, table: str, col: str) -> None:
"""Reset SERIAL sequence to max(col) so future inserts don't collide."""
cur.execute(f"SELECT setval(pg_get_serial_sequence('{table}', '{col}'), COALESCE(MAX({col}), 1)) FROM {table}")
def main() -> None:
print(f"Connecting to PostgreSQL: {DATABASE_URL}")
try:
pg = psycopg2.connect(DATABASE_URL)
except psycopg2.OperationalError as e:
print(f"ERROR: cannot connect to PostgreSQL: {e}", file=sys.stderr)
sys.exit(1)
pg.autocommit = False
print("\n=== economy.sqlite3 ===")
migrate_economy(pg)
print("\n=== penis_stats.sqlite3 (legacy) ===")
migrate_legacy_penis(pg)
print("\n=== bets.sqlite3 ===")
migrate_bets(pg)
print("\n=== chat_history.sqlite3 ===")
migrate_chat(pg)
pg.close()
print("\nMigration complete.")
if __name__ == "__main__":
main()