258 lines
9.0 KiB
Python
258 lines
9.0 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)
|