Files
ccsite-manager/btc_monitor.py

179 lines
6.2 KiB
Python

#!/usr/bin/env python3
"""
BTC Blockchain Monitor — polls blockstream.info for real transactions.
Stores all TXs in SQLite. Detects new deposits. Feeds the dashboard.
"""
import requests
import sqlite3
import time
import json
from datetime import datetime
BTC_ADDRESS = "35vTSjPCaqudcL5cj5ChpNDmqrpPFmEnss"
DB_PATH = "/Users/drjones/ccsite-manager/btc_chain.db"
BLOCKSTREAM_URL = f"https://blockstream.info/api/address/{BTC_ADDRESS}/txs"
def init_db():
conn = sqlite3.connect(DB_PATH)
c = conn.cursor()
c.execute('''CREATE TABLE IF NOT EXISTS transactions (
txid TEXT PRIMARY KEY,
timestamp INTEGER,
amount_btc REAL,
amount_sats INTEGER,
confirmations INTEGER,
fee INTEGER,
raw_json TEXT,
first_seen REAL,
alerted INTEGER DEFAULT 0
)''')
c.execute('CREATE INDEX IF NOT EXISTS idx_tx_time ON transactions(timestamp)')
conn.commit()
conn.close()
def check_blockchain():
"""Poll Blockstream API for transactions. Returns list of new TXs."""
try:
r = requests.get(BLOCKSTREAM_URL, timeout=15)
if r.status_code != 200:
return []
txs = r.json()
except Exception as e:
print(f"Blockstream API error: {e}")
return []
conn = sqlite3.connect(DB_PATH)
c = conn.cursor()
new_txs = []
for tx in txs:
txid = tx['txid']
# Check if we've already seen this TX
c.execute("SELECT txid FROM transactions WHERE txid = ?", (txid,))
if c.fetchone():
continue
# Calculate BTC amount sent TO our address (only count vouts where our address receives)
total_out = 0
for vout in tx.get('vout', []):
addr = vout.get('scriptpubkey_address', '')
if addr == BTC_ADDRESS:
total_out += vout.get('value', 0)
if total_out == 0:
continue # Skip transactions where our address isn't receiving
# Get fee and confirmations
fee = tx.get('fee', 0)
confs = tx.get('status', {}).get('confirmed', False)
conf_count = tx.get('status', {}).get('block_height', 0)
ts = tx.get('status', {}).get('block_time', int(time.time()))
c.execute('''INSERT INTO transactions (txid, timestamp, amount_btc, amount_sats,
confirmations, fee, raw_json, first_seen)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)''',
(txid, ts, total_out, int(total_out * 100000000),
conf_count if confs else 0, fee, json.dumps(tx), time.time()))
new_txs.append({
'txid': txid,
'amount_btc': total_out,
'amount_usd': round(total_out * 68000, 2), # Approximate BTC price
'confirmations': conf_count if confs else 0,
'time': ts,
})
conn.commit()
conn.close()
return new_txs
def auto_complete_orders(new_txs):
"""When new TXs arrive, try to match them to pending orders on all 10 sites."""
SITE_IPS = [f"10.30.20.{ip}" for ip in ['210','211','212','213','214','217','218','219','220','221']]
for tx in new_txs:
btc_amount = tx['amount_btc']
if btc_amount <= 0:
continue
# Try each site
for ip in SITE_IPS:
try:
r = requests.get(f"http://{ip}:5000/api/products", timeout=5)
if r.status_code != 200:
continue
# Check for pending orders via health endpoint
# We need to tell the site "a payment of X BTC arrived, check pending orders"
r2 = requests.post(f"http://{ip}:5000/api/btc-check",
json={"amount_btc": btc_amount, "txid": tx['txid']},
timeout=5)
if r2.status_code == 200:
result = r2.json()
if result.get('matched'):
print(f" ✅ Auto-completed order on {ip}: {result.get('order_code','?')}")
except:
pass
def get_stats():
"""Get aggregate stats from stored transactions."""
conn = sqlite3.connect(DB_PATH)
c = conn.cursor()
c.execute("SELECT COUNT(*), COALESCE(SUM(amount_btc),0), COALESCE(SUM(amount_sats),0) FROM transactions")
count, total_btc, total_sats = c.fetchone()
c.execute("SELECT COUNT(*), COALESCE(SUM(amount_btc),0) FROM transactions WHERE first_seen > ?",
(time.time() - 86400,))
day_count, day_btc = c.fetchone()
c.execute("SELECT txid, amount_btc, timestamp FROM transactions ORDER BY timestamp DESC LIMIT 10")
recent = [{'txid': r[0][:16]+'...', 'btc': r[1], 'ts': r[2]} for r in c.fetchall()]
conn.close()
return {
'total_txs': count,
'total_btc': round(total_btc, 8),
'total_sats': total_sats,
'total_usd': round(total_btc * 68000, 2),
'day_txs': day_count,
'day_btc': round(day_btc, 8),
'day_usd': round(day_btc * 68000, 2),
'recent': recent,
'address': BTC_ADDRESS,
}
def monitor_loop(callback=None):
"""Run continuously, calling callback with new TXs."""
init_db()
print(f"₿ BTC Monitor started — watching {BTC_ADDRESS[:16]}...")
print(f" Polling blockstream.info every 60s")
while True:
try:
new = check_blockchain()
if new:
for tx in new:
print(f" 💰 NEW TX: {tx['amount_btc']:.8f} BTC (${tx['amount_usd']:.2f}) — {tx['txid'][:16]}...")
if callback:
callback(new)
else:
# Silent tick — no new TXs
pass
except Exception as e:
print(f"Monitor error: {e}")
time.sleep(60)
if __name__ == '__main__':
import sys
if len(sys.argv) > 1 and sys.argv[1] == '--once':
new = check_blockchain()
stats = get_stats()
print(json.dumps({'new_txs': len(new), 'stats': stats}, indent=2))
elif len(sys.argv) > 1 and sys.argv[1] == '--stats':
stats = get_stats()
print(json.dumps(stats, indent=2))
else:
monitor_loop()