"""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, VEX_MATURITY_DAYS, VEX_STALE_DAYS) # 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 eligible_voters(owner=None): """V57: pemilih yang memenuhi syarat reward (single source, dipakai `distribute.py` + endpoint `/api/voter/`). Jendela reward V40/V52 = `last_vote > now−STALE` (⊥ basi) DAN akun MATANG (`first_seen_at IS NOT NULL` DAN `first_seen_at ≤ now−MATURITY`). Maturity diukur dari `first_seen_at` (presisi detik; B15). Tanpa `owner` → semua baris `(owner, staked)` utk distribusi; dengan `owner` → baris akun itu (`[]` bila tak memenuhi syarat) utk cek endpoint. """ from datetime import datetime, timedelta, timezone now = datetime.now(timezone.utc).replace(tzinfo=None) stale_cutoff = (now - timedelta(days=VEX_STALE_DAYS)).isoformat(timespec='seconds') mature_cutoff = (now - timedelta(days=VEX_MATURITY_DAYS)).isoformat(timespec='seconds') sql = ('SELECT v.owner, v.staked FROM voters v ' 'LEFT JOIN voter_first_seen fs ON fs.owner = v.owner ' 'WHERE v.last_vote > %s ' ' AND fs.first_seen_at IS NOT NULL ' ' AND fs.first_seen_at <= %s') params = [stale_cutoff, mature_cutoff] if owner is not None: sql += ' AND v.owner = %s' params.append(owner) ensure_first_seen_schema() return query(sql, params) 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'))