Files
databisnisid/db.py
T
proitlab 8e80653df9 klaim baseline saldo stabil (V56/B18) + endpoint cek akun payout-eligible /api/voter (V57)
V56/B18: pre-claim liquid balance dibaca SEKALI per window (poll_claim) dan dibawa ke tiap _try_claim_once — klaim yang mendarat di attempt 1 membuat before attempt 2 (ditolak already-claimed) sudah termasuk reward → delta 0 → alert palsu KLAIM MENDARAT · REWARD TAK TERUKUR (kasus 10 Agt: bfa9ab02 reward 1679.3365 vs tx ditolak b0a8b2f9) ; oracle B18 attempt-1-mendarat → claimed + reward asli terukur

V57: db.py single source eligible_voters(owner=None) — SQL jendela payout V40/V52 (last_vote > now−STALE AND first_seen_at IS NOT NULL AND first_seen_at ≤ now−MATURITY, cutoff UTC-naive) dipakai BAIK distribute.py run_distribution (refactor, hapus _eligible_cutoffs) MAUPUN endpoint baru GET /api/voter/<owner> → 200 {owner, is_valid_voter} (DB-snapshot, selalu 200, ⊖ format-validation); BARU/kadaluarsa/hidden/asing → false ; docs §I/§V/V57 + T50 + AGENTS ; oracle test_dashboard (V57 true/false + drift endpoint==db pool + pool distribusi==db) + test_distribute (mock eligible_voters utk no-op, distribusi luas)
2026-08-12 22:28:42 +07:00

432 lines
16 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""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/<owner>`).
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'))