/** * Backfill of the buyer-pays price on existing sales (ยง35): recomputes all_in_usd / fee_basis / * buyer_premium_rate from the connector flag, the auction house and the fee schedule. Idempotent * and resumable: rows that already carry a fee_basis are skipped unless `force`. Touched assets are * marked for revaluation (asset_stats.updated_at pushed back, like `ri regrade`). */ import { sql } from 'drizzle-orm'; import { logger } from '@rareindex/shared'; import { db } from './lib/db.ts'; import { saleAllIn } from './entity-resolution/writers.ts'; const log = logger.child({ component: 'fees-backfill' }); export interface FeesBackfillResult { scanned: number; updated: number; byBasis: Record; assets: number; } interface Row { id: string; asset_id: string; price_usd: number; currency: string; buyer_premium_included: boolean | null; auction_house: string | null; sale_type: string; source_id: string; connector_id: string; sale_date: Date | string; fee_basis: string | null; } export async function backfillFees(opts: { limit?: number; connectorId?: string; force?: boolean; batch?: number } = {}): Promise { const limit = opts.limit ?? 5_000_000; const batch = opts.batch ?? 5000; const res: FeesBackfillResult = { scanned: 0, updated: 0, byBasis: {}, assets: 0 }; const touched = new Set(); let lastId = ''; while (res.scanned < limit) { const rows = (await db().execute(sql` select s.id, s.asset_id, s.price_usd::float as price_usd, s.currency, s.buyer_premium_included, s.auction_house, s.sale_type, s.source_id, s.connector_id, s.sale_date, s.fee_basis from sales s where s.id > ${lastId} ${opts.connectorId ? sql`and s.connector_id = ${opts.connectorId}` : sql``} ${opts.force ? sql`` : sql`and s.fee_basis is null`} order by s.id limit ${Math.min(batch, limit - res.scanned)}`)) as unknown as Row[]; if (!rows.length) break; res.scanned += rows.length; lastId = rows[rows.length - 1]!.id; const updates: Array<{ id: string; allIn: number; basis: string; rate: number }> = []; for (const r of rows) { const f = await saleAllIn({ priceUsd: Number(r.price_usd), currency: r.currency, buyerPremiumIncluded: r.buyer_premium_included, auctionHouse: r.auction_house, saleType: r.sale_type, sourceId: r.source_id, connectorId: r.connector_id, saleDate: new Date(r.sale_date) }); res.byBasis[f.feeBasis] = (res.byBasis[f.feeBasis] ?? 0) + 1; updates.push({ id: r.id, allIn: f.allInUsd, basis: f.feeBasis, rate: f.buyerPremiumRate }); if (f.feeBasis.startsWith('added_')) touched.add(r.asset_id); } // one statement per batch; the payload travels as JSON (drizzle serialises JS arrays as records, not SQL arrays) await db().execute(sql` update sales s set all_in_usd = u.all_in, fee_basis = u.basis, buyer_premium_rate = u.rate from json_to_recordset(${JSON.stringify(updates.map((u) => ({ id: u.id, all_in: u.allIn, basis: u.basis, rate: u.rate })))}::json) as u(id text, all_in numeric, basis text, rate real) where s.id = u.id`); res.updated += updates.length; log.info({ scanned: res.scanned, updated: res.updated, lastId }, 'fees backfill progress'); } res.assets = touched.size; const ids = [...touched]; for (let i = 0; i < ids.length; i += 1000) await db().execute(sql`update asset_stats set updated_at = '1970-01-01' where asset_id in ${ids.slice(i, i + 1000)}`); return res; }