- Все 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>
369 lines
13 KiB
Python
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()
|