Coverage for app / services / balance_tracker.py: 53%
117 statements
« prev ^ index » next coverage.py v7.13.3, created at 2026-02-04 06:09 -0500
« prev ^ index » next coverage.py v7.13.3, created at 2026-02-04 06:09 -0500
1"""Balance tracking service for updating campaign balances and recording donations."""
2import logging
3from datetime import datetime, timezone
4from typing import List, Dict, Any, Optional
5from sqlalchemy.ext.asyncio import AsyncSession
6from sqlalchemy import select
8from app.db.models import Campaign, Donation, CampaignStatus
9from app.services.blockchain import BlockchainService, BlockchainAPIError
11logger = logging.getLogger(__name__)
14class BalanceTracker:
15 """Service for tracking campaign balances and recording donations."""
17 def __init__(self, db: AsyncSession):
18 self.db = db
19 self.blockchain_service = BlockchainService()
21 async def update_campaign_balances(self, campaign_id: str) -> None:
22 """Update balances for a single campaign."""
23 result = await self.db.execute(
24 select(Campaign).where(Campaign.id == campaign_id)
25 )
26 campaign = result.scalar_one_or_none()
28 if not campaign:
29 return
31 # Update BTC balance
32 if campaign.btc_wallet_address:
33 try:
34 balance = await self.blockchain_service.get_btc_balance(campaign.btc_wallet_address)
35 campaign.current_btc_satoshi = balance
36 except BlockchainAPIError as e:
37 logger.warning(f"Failed to update BTC balance for campaign {campaign_id}: {e}")
39 # Update ETH balance
40 if campaign.eth_wallet_address:
41 try:
42 balance = await self.blockchain_service.get_eth_balance(campaign.eth_wallet_address)
43 campaign.current_eth_wei = balance
44 except BlockchainAPIError as e:
45 logger.warning(f"Failed to update ETH balance for campaign {campaign_id}: {e}")
47 # Update DOGE balance
48 if campaign.doge_wallet_address:
49 try:
50 balance = await self.blockchain_service.get_doge_balance(campaign.doge_wallet_address)
51 campaign.current_doge_satoshi = balance
52 except BlockchainAPIError as e:
53 logger.warning(f"Failed to update DOGE balance for campaign {campaign_id}: {e}")
55 # Update SOL balance
56 if campaign.sol_wallet_address:
57 try:
58 balance = await self.blockchain_service.get_sol_balance(campaign.sol_wallet_address)
59 campaign.current_sol_lamports = balance
60 except BlockchainAPIError as e:
61 logger.warning(f"Failed to update SOL balance for campaign {campaign_id}: {e}")
63 campaign.last_balance_check = datetime.now(timezone.utc)
64 await self.db.commit()
66 async def record_new_donations(
67 self,
68 campaign_id: str,
69 chain: str,
70 transactions: List[Dict[str, Any]]
71 ) -> int:
72 """Record new donations from transactions. Returns count of new donations."""
73 # Fetch campaign to get wallet addresses
74 campaign_result = await self.db.execute(
75 select(Campaign).where(Campaign.id == campaign_id)
76 )
77 campaign = campaign_result.scalar_one_or_none()
79 if not campaign:
80 return 0
82 # Get existing donation tx hashes to avoid duplicates
83 result = await self.db.execute(
84 select(Donation.tx_hash).where(
85 Donation.campaign_id == campaign_id,
86 Donation.chain == chain
87 )
88 )
89 existing_hashes = set(result.scalars().all())
91 new_donations = []
92 for tx in transactions:
93 tx_hash = tx.get("hash") or tx.get("tx_hash") or tx.get("signature")
94 if not tx_hash or tx_hash in existing_hashes:
95 continue
97 # Extract amount and from_address based on chain
98 amount = 0
99 from_address = None
101 if chain == "btc":
102 # BlockCypher BTC format
103 outputs = tx.get("outputs", [])
104 for output in outputs:
105 addresses = output.get("addresses", [])
106 if campaign.btc_wallet_address in addresses:
107 amount += output.get("value", 0)
108 inputs = tx.get("inputs", [])
109 if inputs:
110 from_address = inputs[0].get("addresses", [None])[0]
112 elif chain == "eth":
113 # BlockCypher ETH format
114 outputs = tx.get("outputs", [])
115 for output in outputs:
116 addresses = output.get("addresses", [])
117 if campaign.eth_wallet_address in addresses:
118 amount += output.get("value", 0)
119 inputs = tx.get("inputs", [])
120 if inputs:
121 from_address = inputs[0].get("addresses", [None])[0]
123 elif chain == "doge":
124 # Same as BTC
125 outputs = tx.get("outputs", [])
126 for output in outputs:
127 addresses = output.get("addresses", [])
128 if campaign.doge_wallet_address in addresses:
129 amount += output.get("value", 0)
130 inputs = tx.get("inputs", [])
131 if inputs:
132 from_address = inputs[0].get("addresses", [None])[0]
134 elif chain == "sol":
135 # Helius SOL format
136 native_transfers = tx.get("nativeTransfers", [])
137 for transfer in native_transfers:
138 if transfer.get("toUserAccount") == campaign.sol_wallet_address:
139 amount += transfer.get("amount", 0)
140 from_address = tx.get("fromUserAccount")
142 if amount > 0:
143 donation = Donation(
144 campaign_id=campaign_id,
145 chain=chain,
146 tx_hash=tx_hash,
147 amount_smallest_unit=amount,
148 from_address=from_address,
149 confirmed_at=datetime.now(timezone.utc),
150 block_number=tx.get("block_height") or tx.get("slot"),
151 )
152 new_donations.append(donation)
154 if new_donations:
155 self.db.add_all(new_donations)
156 await self.db.commit()
158 return len(new_donations)
160 async def poll_all_active_campaigns(self) -> None:
161 """Poll balances for all active campaigns."""
162 result = await self.db.execute(
163 select(Campaign).where(Campaign.status == CampaignStatus.ACTIVE)
164 )
165 campaigns = result.scalars().all()
167 for campaign in campaigns:
168 try:
169 await self.update_campaign_balances(campaign.id)
171 # Also check for new transactions and record donations
172 if campaign.btc_wallet_address:
173 txs = await self.blockchain_service.get_btc_transactions(campaign.btc_wallet_address, limit=10)
174 await self.record_new_donations(campaign.id, "btc", txs)
176 if campaign.eth_wallet_address:
177 txs = await self.blockchain_service.get_eth_transactions(campaign.eth_wallet_address, limit=10)
178 await self.record_new_donations(campaign.id, "eth", txs)
180 if campaign.doge_wallet_address:
181 txs = await self.blockchain_service.get_doge_transactions(campaign.doge_wallet_address, limit=10)
182 await self.record_new_donations(campaign.id, "doge", txs)
184 if campaign.sol_wallet_address:
185 txs = await self.blockchain_service.get_sol_transactions(campaign.sol_wallet_address, limit=10)
186 await self.record_new_donations(campaign.id, "sol", txs)
187 except Exception as e:
188 logger.error(f"Error polling campaign {campaign.id}: {e}", exc_info=True)