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

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 

7 

8from app.db.models import Campaign, Donation, CampaignStatus 

9from app.services.blockchain import BlockchainService, BlockchainAPIError 

10 

11logger = logging.getLogger(__name__) 

12 

13 

14class BalanceTracker: 

15 """Service for tracking campaign balances and recording donations.""" 

16 

17 def __init__(self, db: AsyncSession): 

18 self.db = db 

19 self.blockchain_service = BlockchainService() 

20 

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() 

27 

28 if not campaign: 

29 return 

30 

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}") 

38 

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}") 

46 

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}") 

54 

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}") 

62 

63 campaign.last_balance_check = datetime.now(timezone.utc) 

64 await self.db.commit() 

65 

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() 

78 

79 if not campaign: 

80 return 0 

81 

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()) 

90 

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 

96 

97 # Extract amount and from_address based on chain 

98 amount = 0 

99 from_address = None 

100 

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] 

111 

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] 

122 

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] 

133 

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") 

141 

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) 

153 

154 if new_donations: 

155 self.db.add_all(new_donations) 

156 await self.db.commit() 

157 

158 return len(new_donations) 

159 

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() 

166 

167 for campaign in campaigns: 

168 try: 

169 await self.update_campaign_balances(campaign.id) 

170 

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) 

175 

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) 

179 

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) 

183 

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)