Files
proitlab 1e93437a90 feat: VOTE ULANG SEGERA Telegram notification + emoticon on all notifications
- telegram.py: emoticons on all existing notif (🚀✅⚠️🔴💰🔍) + new notify_warn for expiring voters
- db.py: voter_notify table (owner/urgency/notified_at) for state-change dedup
- get_voters.py: _notify_warn() after persist() — urgency escalation ⚠️/🔴/🚨 by days remaining
- Dockerfile.scan: add telegram.py to COPY
- test_get_voters.py: 15 new V66 assertions (schema, record, cleanup, integration, escalation)
- test_distribute.py: updated notification format assertions
- SPEC.md: V66 + B26
2026-08-27 16:12:53 +07:00

520 lines
19 KiB
Python
Raw Permalink 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, timezone
if updated_at is None:
updated_at = datetime.now(timezone.utc).replace(tzinfo=None).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, timezone
conn = connect()
try:
cur = conn.cursor()
if fee_status == 'sent':
now = datetime.now(timezone.utc).replace(tzinfo=None).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()
# ——————————————————————————————————————————————————————————————————————
# V66: Notifikasi VOTE ULANG SEGERA — tabel lacak siapa sudah diberi tahu
# pada level urgensi berapa. Satu baris per owner (PK), di-update saat
# urgensi naik. Bersihkan otomatis saat pemilih keluar dari jendela warn.
# ——————————————————————————————————————————————————————————————————————
_NOTIFY_COLS = (
'owner VARCHAR(13) PRIMARY KEY, '
'urgency INTEGER NOT NULL, '
'notified_at TEXT NOT NULL'
)
_NOTIFY_DDL_SQLITE = (f'CREATE TABLE IF NOT EXISTS voter_notify '
f'({_NOTIFY_COLS})')
_NOTIFY_DDL_MYSQL = (f'CREATE TABLE IF NOT EXISTS voter_notify '
f'({_NOTIFY_COLS})')
def ensure_notify_schema():
"""Buat tabel lacak notifikasi VOTE ULANG SEGERA bila belum ada (idempoten)."""
conn = connect()
try:
cur = conn.cursor()
if DB_BACKEND == 'sqlite':
cur.execute(_NOTIFY_DDL_SQLITE)
else:
cur.execute(_NOTIFY_DDL_MYSQL)
conn.commit()
finally:
conn.close()
def get_notify_urgency(owner):
"""V66: ambil level urgensi terakhir yang sudah diberi tahu utk `owner`.
Kembalikan integer (1/2/3) atau None bila belum pernah diberi tahu.
"""
ensure_notify_schema()
row = queryone(
'SELECT urgency FROM voter_notify WHERE owner = %s', (owner,))
return row[0] if row else None
def record_notify(owner, urgency, notified_at):
"""V66: catat bahwa `owner` sudah diberi tahu pada level `urgency`.
INSERT OR IGNORE (baru) + UPDATE bila urgensi naik (tidak turun).
Idempoten — safe dipanggil berulang.
"""
ensure_notify_schema()
conn = connect()
try:
cur = conn.cursor()
if DB_BACKEND == 'sqlite':
cur.execute(
'INSERT OR IGNORE INTO voter_notify '
'(owner, urgency, notified_at) VALUES (?, ?, ?)',
(owner, urgency, notified_at))
cur.execute(
'UPDATE voter_notify SET urgency = ?, notified_at = ? '
'WHERE owner = ? AND urgency < ?',
(urgency, notified_at, owner, urgency))
else:
cur.execute(
'INSERT INTO voter_notify '
'(owner, urgency, notified_at) VALUES (%s, %s, %s) '
'ON DUPLICATE KEY UPDATE '
'urgency = IF(%s > urgency, %s, urgency), '
'notified_at = IF(%s > urgency, %s, notified_at)',
(owner, urgency, notified_at, urgency, urgency,
urgency, notified_at))
conn.commit()
finally:
conn.close()
def cleanup_notify(cutoff):
"""V66: hapus baris notifikasi yang sudah lewat jendela warn (stale)."""
ensure_notify_schema()
conn = connect()
try:
cur = conn.cursor()
cur.execute(_translate(
'DELETE FROM voter_notify WHERE notified_at < %s'), (cutoff,))
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'))