🌳
pt0/donateapp/serverF/paymentScanAI.mts
1import type { Knex } from 'knex'
2import * as _ from 'lodash-es'
3import { randomBytes } from 'crypto'
8import { toEthAddrLower, type EthAddrLower } from './ethAddrAI.mts'
11const donateNetworkNames = ['homestead', 'arbitrum', 'optimism', 'base', 'linea'] as const
13interface PendingEntry {
14 id: string
15 dona_ref: string
16 dona_ref_kind: string
17 support_oppose: string
18 amount_usd: string
19 eth_gwei: string
20 eth_payment_addr: string
21 expires_at: Date
22 created_at: Date
23 charity_name: string
24 target_entry_id: string | null
27const ethScanTxLock = async ({db}: {db: Knex | Knex.Transaction}) => {
28 await db.raw('SELECT pg_advisory_xact_lock(760491);')
31export const scanDonaOrderPayment = async ({db, order}: {db: Knex, order: PendingEntry}) => {
32 const ethPaymentAddr = order.eth_payment_addr
33 const targetGwei = parseInt(order.eth_gwei, 10)
35 let paid = false
36 let entryId: string | null = null
37 let uploadToken: string | null = null
39 // Sequential to avoid pool exhaustion - parallel transactions all wait on advisory lock 760491.
40 // If a transaction hangs (e.g. RPC timeout), idle_in_transaction_session_timeout (60s, pg config) kills it.
41 for (const networkName of donateNetworkNames) {
42 if (paid) break
44 const cacheNetworkName = networkName === 'homestead' ? null : networkName
45 const key = _.compact(['startblock', cacheNetworkName, ethPaymentAddr]).join('-')
47 await db.transaction(async (trx: Knex.Transaction) => {
48 await ethScanTxLock({db: trx})
49 const [getStartBlock, setStartBlock] = genericKvGetSet({db: trx, key, defaultValue: 0})
50 const startBlock = Number(await getStartBlock()) || 0
52 const transactions = await cacheFetch({
53 cacheKey: `scanMillionpx-${key}`,
54 ttlMs: 5 * 1000,
55 cacheMissFn: async () => {
56 return await getEtherTxs({address: ethPaymentAddr, startBlock, networkName})
57 }
58 })
60 const newBlockNo = _.max(_.map(transactions, 'blockNumber').map(Number))
62 const orderCreatedAt = Math.floor(order.created_at.getTime() / 1000)
63 for (const paymentTrx of transactions) {
64 const {eth_gwei, hash: eth_txid, from: senderAddr, timestamp: txTimestamp} = paymentTrx
65 if (txTimestamp < orderCreatedAt) continue
67 if (eth_gwei === targetGwei) {
68 const result = await finalizeDonaPayment({db: trx, order, eth_txid, senderAddr})
69 paid = true
70 entryId = result.entryId
71 uploadToken = result.uploadToken
72 // Only advance startBlock after successful payment detection
73 // This prevents skipping over txs that Alchemy hasn't indexed yet
74 if (newBlockNo != null && newBlockNo > startBlock) {
75 await setStartBlock(newBlockNo)
76 }
77 break
78 }
79 }
80 })
81 }
83 return {paid, entryId, uploadToken}
86export const finalizeDonaPayment = async ({db, order, eth_txid, senderAddr: rawSenderAddr}: {db: Knex | Knex.Transaction, order: PendingEntry, eth_txid: string, senderAddr: string}) => {
87 const senderAddr: EthAddrLower = toEthAddrLower(rawSenderAddr)
88 const amountUsd = parseFloat(order.amount_usd)
89 let entryId: string
91 if (order.target_entry_id) {
92 // Chip-in: add funds to existing entry, delete the pending order entry
93 await db('dona_entries')
94 .where({id: order.target_entry_id})
95 .increment('amount_usd', amountUsd)
96 await db('dona_entries').where({id: order.id}).delete()
97 entryId = order.target_entry_id
98 } else {
99 // New entry: check capacity, then finalize this entry
100 const countResult = await paidEntries(db).count('* as count').first()
101 const filledCount = parseInt(countResult?.count as string || '0', 10)
103 if (filledCount >= maxDonaEntries) {
104 const lowestEntry = await paidEntries(db).orderBy('amount_usd', 'asc').first()
105 if (lowestEntry && amountUsd > parseFloat(lowestEntry.amount_usd)) {
106 await db('dona_entries').where({id: lowestEntry.id}).delete()
107 }
108 }
110 // Finalize the pending entry by setting payment fields and clearing order fields
111 const uploadToken = order.dona_ref_kind === 'ipfs_cid' ? randomBytes(32).toString('hex') : null
112 await db('dona_entries').where({id: order.id}).update({
113 eth_txid,
114 sender_addr: senderAddr,
115 paid_at: db.fn.now(),
116 eth_gwei: null,
117 eth_payment_addr: null,
118 expires_at: null,
119 target_entry_id: null,
120 upload_token: uploadToken,
121 })
122 entryId = order.id
124 // Record payment for charity tracking
125 await db('dona_payments').insert({
126 entry_id: entryId,
127 amount_usd: amountUsd,
128 charity_name: order.charity_name,
129 paid_at: db.fn.now(),
130 source_order_id: order.id,
131 sender_addr: senderAddr,
132 })
134 return {entryId, uploadToken}
135 }
137 // Chip-in: record payment
138 await db('dona_payments').insert({
139 entry_id: entryId,
140 amount_usd: amountUsd,
141 charity_name: order.charity_name,
142 paid_at: db.fn.now(),
143 source_order_id: order.id,
144 sender_addr: senderAddr,
145 })
147 return {entryId, uploadToken: null}