Files
databisnisid/db.py
T

405 lines
15 KiB
Python

"""Abstraksi penyimpanan: sqlite (default) atau MariaDB/MySQL (PyMySQL).
Pilih backend via config `VEX_DB_BACKEND` (env `.env`). Kedua backend memakai
skema `voters` yang sama; query berbagi sintaks ditulis dengan placeholder
`%s` yang diterjemahkan ke `?` untuk sqlite. Dashboard read-only (V10).
"""
import sqlite3
from config import (DB_BACKEND, DB_HOST, DB_NAME, DB_PASS, DB_PORT, DB_PATH,
DB_USER)
# Skema voters — sama untuk kedua backend (I.db)
_COLS = (
'owner VARCHAR(13) PRIMARY KEY, '
'weight TEXT NOT NULL, '
'staked DOUBLE NOT NULL, '
'scanned_at TEXT NOT NULL, '
'last_vote TEXT'
)
_SCHEMA = f'CREATE TABLE voters ({_COLS})'
_SCHEMA_IF_NOT_EXISTS = f'CREATE TABLE IF NOT EXISTS voters ({_COLS})'
def _translate(sql):
"""sqlite pakai `?`; mysql pakai `%s`. Ubah `%s` → `?` untuk sqlite."""
return sql.replace('%s', '?') if DB_BACKEND == 'sqlite' else sql
def connect():
"""Buka koneksi backend. sqlite: WAL agar pembaca tak terblokir."""
if DB_BACKEND == 'sqlite':
conn = sqlite3.connect(DB_PATH)
conn.execute('PRAGMA journal_mode=WAL')
return conn
if DB_BACKEND == 'mysql':
import pymysql
return pymysql.connect(
host=DB_HOST, port=DB_PORT, user=DB_USER, password=DB_PASS,
database=DB_NAME, charset='utf8mb4',
)
raise ValueError(f'VEX_DB_BACKEND tak dikenal: {DB_BACKEND!r}')
def query(sql, params=()):
"""Jalankan SELECT → daftar baris (tuple), koneksi dibuka-tutup tiap panggil."""
conn = connect()
try:
cur = conn.cursor()
cur.execute(_translate(sql), params)
return list(cur.fetchall()) # sqlite→list, PyMySQL→tuple: seragamkan
finally:
conn.close()
def queryone(sql, params=()):
"""Jalankan SELECT → baris pertama (tuple) atau None."""
conn = connect()
try:
cur = conn.cursor()
cur.execute(_translate(sql), params)
return cur.fetchone()
finally:
conn.close()
def replace_snapshot(rows, scanned_at):
"""V7: ganti snapshot voters tiap skan (⊥ history/append).
sqlite: DROP+CREATE+INSERT (WAL melindungi pembaca).
mysql: CREATE IF NOT EXISTS + DELETE + INSERT dalam satu transaksi
→ MVCC memberi pembaca snapshot konsisten (⊥ torn read saat ganti harian).
"""
tuples = [(r['owner'], r['weight'], r['staked'], scanned_at, r.get('last_vote'))
for r in rows]
conn = connect()
try:
cur = conn.cursor()
if DB_BACKEND == 'sqlite':
cur.execute('DROP TABLE IF EXISTS voters')
cur.execute(_SCHEMA)
cur.executemany(
'INSERT INTO voters (owner, weight, staked, scanned_at, last_vote) '
'VALUES (?, ?, ?, ?, ?)',
tuples,
)
else:
cur.execute(_SCHEMA_IF_NOT_EXISTS)
cur.execute('DELETE FROM voters')
cur.executemany(
'INSERT INTO voters (owner, weight, staked, scanned_at, last_vote) '
'VALUES (%s, %s, %s, %s, %s)',
tuples,
)
conn.commit()
finally:
conn.close()
# ——————————————————————————————————————————————————————————————————————
# Riwayat distribusi profit-share (V19..V27). Tabel append — tiap transfer
# & tiap run tercatat permanen (⊥ hapus). Skema sama untuk kedua backend,
# beda hanya klausa primary key autoincrement.
# ——————————————————————————————————————————————————————————————————————
_RUN_COLS = (
'run_date TEXT NOT NULL, '
'balance_start DOUBLE NOT NULL, '
'total_voters INTEGER NOT NULL, '
'total_staked DOUBLE NOT NULL, '
'total_sent DOUBLE NOT NULL, '
'status TEXT NOT NULL, '
'created_at TEXT NOT NULL'
)
_PAY_COLS = (
'run_id INTEGER NOT NULL, '
'owner VARCHAR(13) NOT NULL, '
'amount DOUBLE NOT NULL, '
'txid TEXT, '
'status TEXT NOT NULL, '
'error TEXT, '
'created_at TEXT NOT NULL, '
'updated_at TEXT NOT NULL'
)
_RUN_DDL_SQLITE = f'CREATE TABLE IF NOT EXISTS distribute_runs (run_id INTEGER PRIMARY KEY AUTOINCREMENT, {_RUN_COLS})'
_PAY_DDL_SQLITE = f'CREATE TABLE IF NOT EXISTS distribute_payments (payment_id INTEGER PRIMARY KEY AUTOINCREMENT, {_PAY_COLS})'
_RUN_DDL_MYSQL = f'CREATE TABLE IF NOT EXISTS distribute_runs (run_id INT AUTO_INCREMENT PRIMARY KEY, {_RUN_COLS})'
_PAY_DDL_MYSQL = f'CREATE TABLE IF NOT EXISTS distribute_payments (payment_id INT AUTO_INCREMENT PRIMARY KEY, {_PAY_COLS})'
def ensure_distribute_schema():
"""Buat tabel riwayat distribusi bila belum ada (idempoten, kedua backend)."""
conn = connect()
try:
cur = conn.cursor()
if DB_BACKEND == 'sqlite':
cur.execute(_RUN_DDL_SQLITE)
cur.execute(_PAY_DDL_SQLITE)
else:
cur.execute(_RUN_DDL_MYSQL)
cur.execute(_PAY_DDL_MYSQL)
_drop_pay_memo(cur)
conn.commit()
finally:
conn.close()
def _drop_pay_memo(cur):
"""Migrasi: buang kolom memo lama dari distribute_payments (idempoten).
Kolom `memo` pernah disimpan (V23) lalu dihapus — DB produksi yang sudah
punya kolom perlu di-alter; DB baru langsung tanpa kolom.
"""
if DB_BACKEND == 'sqlite':
cur.execute(
"SELECT name FROM pragma_table_info('distribute_payments')")
cols = {row[0] for row in cur.fetchall()}
else:
cur.execute(
"SELECT COLUMN_NAME FROM information_schema.COLUMNS "
"WHERE TABLE_SCHEMA = DATABASE() "
" AND TABLE_NAME = 'distribute_payments'")
cols = {row[0] for row in cur.fetchall()}
if 'memo' in cols:
cur.execute(_translate('ALTER TABLE distribute_payments DROP COLUMN memo'))
def record_run(run_date, balance_start, total_voters, total_staked,
total_sent, status, created_at):
"""Insert baris `distribute_runs`, kembalikan run_id."""
ensure_distribute_schema()
conn = connect()
try:
cur = conn.cursor()
cur.execute(_translate(
'INSERT INTO distribute_runs '
'(run_date, balance_start, total_voters, total_staked, total_sent, '
' status, created_at) VALUES (%s, %s, %s, %s, %s, %s, %s)'),
(run_date, balance_start, total_voters, total_staked, total_sent,
status, created_at))
conn.commit()
return cur.lastrowid
finally:
conn.close()
def record_payment(run_id, owner, amount, status, created_at,
txid=None, error=None):
"""Insert baris `distribute_payments` → payment_id."""
ensure_distribute_schema()
conn = connect()
try:
cur = conn.cursor()
cur.execute(_translate(
'INSERT INTO distribute_payments '
'(run_id, owner, amount, txid, status, error, created_at, '
' updated_at) VALUES (%s, %s, %s, %s, %s, %s, %s, %s)'),
(run_id, owner, amount, txid, status, error, created_at,
created_at))
conn.commit()
return cur.lastrowid
finally:
conn.close()
def update_payment_status(payment_id, status, txid=None, error=None, updated_at=None):
"""Perbarui status satu transfer (V21): pending → sent|failed."""
from datetime import datetime
if updated_at is None:
updated_at = datetime.now().isoformat(timespec='seconds')
conn = connect()
try:
cur = conn.cursor()
cur.execute(_translate(
'UPDATE distribute_payments SET status = %s, txid = %s, '
'error = %s, updated_at = %s WHERE payment_id = %s'),
(status, txid, error, updated_at, payment_id))
conn.commit()
finally:
conn.close()
def update_run_status(run_id, status, total_sent):
"""Tutup run: status ok|partial + total terkirim nyata."""
conn = connect()
try:
cur = conn.cursor()
cur.execute(_translate(
'UPDATE distribute_runs SET status = %s, total_sent = %s '
'WHERE run_id = %s'),
(status, total_sent, run_id))
conn.commit()
finally:
conn.close()
def list_runs():
"""Ringkasan run distribusi, terbaru dulu."""
ensure_distribute_schema()
return query(
'SELECT run_id, run_date, balance_start, total_voters, total_staked, '
' total_sent, status, created_at '
'FROM distribute_runs ORDER BY run_id DESC')
def list_payments(run_id=None):
"""Transfer per-pemilih, terbaru dulu; opsional filter satu run."""
ensure_distribute_schema()
sql = ('SELECT payment_id, run_id, owner, amount, txid, status, error, '
' created_at, updated_at FROM distribute_payments')
params = ()
if run_id is not None:
sql += ' WHERE run_id = %s'
params = (run_id,)
sql += ' ORDER BY payment_id DESC'
return query(sql, params)
# ——————————————————————————————————————————————————————————————————————
# Riwayat klaim reward BP (V32..V34). Satu baris per klaim sukses; fee yang
# belum terkirim (`pending`/`failed`) jadi sumber resume saat restart.
# ——————————————————————————————————————————————————————————————————————
_CLAIM_COLS = (
'run_date TEXT NOT NULL, '
'claim_txid TEXT, '
'reward DOUBLE NOT NULL, '
'fee_amount DOUBLE NOT NULL, '
'fee_status TEXT NOT NULL, ' # pending | sent | failed | skipped
'fee_txid TEXT, '
'claimed_at TEXT, '
'fee_sent_at TEXT, '
'created_at TEXT NOT NULL'
)
_CLAIM_DDL_SQLITE = (f'CREATE TABLE IF NOT EXISTS claim_runs '
f'(claim_id INTEGER PRIMARY KEY AUTOINCREMENT, '
f'{_CLAIM_COLS})')
_CLAIM_DDL_MYSQL = (f'CREATE TABLE IF NOT EXISTS claim_runs '
f'(claim_id INT AUTO_INCREMENT PRIMARY KEY, '
f'{_CLAIM_COLS})')
def ensure_claim_schema():
"""Buat tabel riwayat klaim bila belum ada (idempoten, kedua backend)."""
conn = connect()
try:
cur = conn.cursor()
if DB_BACKEND == 'sqlite':
cur.execute(_CLAIM_DDL_SQLITE)
else:
cur.execute(_CLAIM_DDL_MYSQL)
conn.commit()
finally:
conn.close()
def record_claim(run_date, claim_txid, reward, fee_amount, fee_status,
created_at):
"""Insert baris `claim_runs` → claim_id."""
ensure_claim_schema()
conn = connect()
try:
cur = conn.cursor()
cur.execute(_translate(
'INSERT INTO claim_runs '
'(run_date, claim_txid, reward, fee_amount, fee_status, '
' claimed_at, fee_sent_at, created_at) '
'VALUES (%s, %s, %s, %s, %s, %s, %s, %s)'),
(run_date, claim_txid, reward, fee_amount, fee_status, created_at,
None, created_at))
conn.commit()
return cur.lastrowid
finally:
conn.close()
def update_claim_fee(claim_id, fee_status, fee_txid=None):
"""Perbarui status fee satu klaim (V33): pending → sent|failed|skipped."""
from datetime import datetime
conn = connect()
try:
cur = conn.cursor()
if fee_status == 'sent':
now = datetime.now().isoformat(timespec='seconds')
cur.execute(_translate(
'UPDATE claim_runs SET fee_status = %s, fee_txid = %s, '
'fee_sent_at = %s WHERE claim_id = %s'),
(fee_status, fee_txid, now, claim_id))
else:
cur.execute(_translate(
'UPDATE claim_runs SET fee_status = %s, fee_txid = %s '
'WHERE claim_id = %s'),
(fee_status, fee_txid, claim_id))
conn.commit()
finally:
conn.close()
# ——————————————————————————————————————————————————————————————————————
# Kemunculan pertama pemilih (V41). Tabel append kecil — satu baris per owner
# (PK), mengingat KAPAN pertama kali muncul di snapshot. `voters` tetap
# snapshot murni (diganti tiap skan); first_seen bertahan di sini. Dipakai
# dashboard utk membedakan PEMILIH BARU (belum pernah terlihat) vs REVOTE
# (pemilih yang sudah pernah ada, revote/perpanjang).
# ——————————————————————————————————————————————————————————————————————
_FIRST_SEEN_COLS = 'owner VARCHAR(13) PRIMARY KEY, first_seen_at TEXT NOT NULL'
_FIRST_SEEN_DDL_SQLITE = (f'CREATE TABLE IF NOT EXISTS voter_first_seen '
f'({_FIRST_SEEN_COLS})')
_FIRST_SEEN_DDL_MYSQL = (f'CREATE TABLE IF NOT EXISTS voter_first_seen '
f'({_FIRST_SEEN_COLS})')
def ensure_first_seen_schema():
"""Buat tabel kemunculan pertama bila belum ada (idempoten, kedua backend)."""
conn = connect()
try:
cur = conn.cursor()
if DB_BACKEND == 'sqlite':
cur.execute(_FIRST_SEEN_DDL_SQLITE)
else:
cur.execute(_FIRST_SEEN_DDL_MYSQL)
conn.commit()
finally:
conn.close()
def record_first_seen(rows, scanned_at):
"""Catat kemunculan pertama tiap owner (idempoten, ⊥ overwrite).
INSERT OR IGNORE (sqlite) / INSERT IGNORE (mysql): hanya baris yang belum
ada yang masuk — `first_seen_at` selalu tanggal kemunculan pertama.
"""
ensure_first_seen_schema()
conn = connect()
try:
cur = conn.cursor()
if DB_BACKEND == 'sqlite':
sql = ('INSERT OR IGNORE INTO voter_first_seen '
'(owner, first_seen_at) VALUES (?, ?)')
else:
sql = ('INSERT IGNORE INTO voter_first_seen '
'(owner, first_seen_at) VALUES (%s, %s)')
cur.executemany(sql, [(r['owner'], scanned_at) for r in rows])
conn.commit()
finally:
conn.close()
def pending_claim_fee():
"""Klaim terbaru dengan fee belum terkirim (fee_amount > 0) → baris | None.
Source of truth resume: fee `pending`/`failed` diulang tiap siklus dan
dilanjutkan saat restart (V34).
"""
ensure_claim_schema()
return queryone(
'SELECT claim_id, run_date, claim_txid, reward, fee_amount, '
' fee_status, fee_txid, claimed_at, fee_sent_at, created_at '
'FROM claim_runs '
'WHERE fee_status IN (%s, %s) AND fee_amount > 0 '
'ORDER BY claim_id DESC LIMIT 1',
('pending', 'failed'))