import { and, desc, eq, gte, inArray, isNull, or, sql } from 'drizzle-orm'; import { assetStats, assetVariants, assets, gradePremiums, listings, populationReports, priceObservations, priceSnapshots, sales, sources, valuations, variantStats } from '@rareindex/database'; import { ASK_ANOMALY_HIGH_RATIO, ASK_ANOMALY_LOW_RATIO, ASK_MIN_CONFIDENCE, ASK_MIN_MATCH_CONFIDENCE, ASK_MIN_SAMPLE, DEAL_REVIEW_THRESHOLD, adjustmentFactor, assetDataQuality, computeGradePremiums, computeValuation, detectOutliers, liquidityScore, momentumScore, rarityScore, trendingScore, type GradePremium, type SaleInput } from '@rareindex/valuation'; import { logger, median, newId, pctChange, toDateOnly } from '@rareindex/shared'; import { db } from '../lib/db.ts'; import { auditMany } from '../lib/audit.ts'; import { emit } from '../lib/events.ts'; const log = logger.child({ component: 'valuation' }); const DAY = 86_400_000; const safeRatioV = (v: number | null): number | null => (v === null || !Number.isFinite(v) || Math.abs(v) >= 99 ? null : v); const asDate = (v: unknown): Date | null => (v instanceof Date ? v : typeof v === 'string' ? new Date(v) : null); const premiumCache = new Map(); async function premiumsFor(categorySlug: string): Promise { const hit = premiumCache.get(categorySlug); if (hit && Date.now() - hit.at < 3600_000) return hit.rows; const rows = await db().select({ grader: gradePremiums.grader, grade: gradePremiums.grade, marketMultiplier: gradePremiums.marketMultiplier, sampleSize: gradePremiums.sampleSize }).from(gradePremiums).where(eq(gradePremiums.categorySlug, categorySlug)); premiumCache.set(categorySlug, { at: Date.now(), rows }); return rows; } const trustCache = new Map(); async function trustMap(): Promise> { if (trustCache.size) return trustCache; for (const s of await db().select({ id: sources.id, t: sources.trustScore }).from(sources)) trustCache.set(s.id, Number(s.t)); return trustCache; } interface SaleRow { id: string; variantId: string | null; sourceId: string; saleDate: Date; /** buyer-pays price: coalesce(all_in_usd, price_usd) — hammer + estimated premium when applicable (§35) */ priceUsd: number; /** included | added_published | added_approximate | added_default | none | unknown | null */ feeBasis: string | null; grader: string | null; grade: string | null; quantity: number; isBundle: boolean; status: string; confidence: number; flags: string[]; } const BUYER_PAYS = sql`coalesce(${sales.allInUsd}, ${sales.priceUsd})`; /** Value one asset: outliers → per-variant valuations → asset stats, snapshots, listing discounts. */ export async function valueAsset(assetId: string, opts: { now?: Date; rebuildHistory?: boolean } = {}): Promise<{ variants: number; riv: number | null }> { const now = opts.now ?? new Date(); const [asset] = await db().select().from(assets).where(eq(assets.id, assetId)).limit(1); if (!asset) return { variants: 0, riv: null }; const trust = await trustMap(); const variants = await db().select().from(assetVariants).where(eq(assetVariants.assetId, assetId)); const since3y = new Date(now.getTime() - 3 * 365 * DAY); const saleRows = (await db() .select({ id: sales.id, variantId: sales.variantId, sourceId: sales.sourceId, saleDate: sales.saleDate, priceUsd: BUYER_PAYS, feeBasis: sales.feeBasis, grader: sales.grader, grade: sales.grade, quantity: sales.quantity, isBundle: sales.isBundle, status: sales.status, confidence: sales.confidence, flags: sales.flags }) .from(sales) .where(and(eq(sales.assetId, assetId), gte(sales.saleDate, since3y))) .orderBy(desc(sales.saleDate))) as SaleRow[]; const allSalesCount = (await db().select({ n: sql`count(*)::int`, min: sql`min(${BUYER_PAYS})`, max: sql`max(${BUYER_PAYS})`, minAt: sql`(array_agg(${sales.saleDate} order by ${BUYER_PAYS} asc))[1]`, maxAt: sql`(array_agg(${sales.saleDate} order by ${BUYER_PAYS} desc))[1]` }).from(sales).where(and(eq(sales.assetId, assetId), eq(sales.status, 'valid'))))[0]!; // 1. outliers per variant (flag only new ones; never delete) const flagsToApply: Array<{ id: string; reason: string; score: number }> = []; const byVariant = new Map(); for (const s of saleRows) byVariant.set(s.variantId ?? '', [...(byVariant.get(s.variantId ?? '') ?? []), s]); for (const list of byVariant.values()) { const valid = list.filter((s) => s.status === 'valid'); for (const f of detectOutliers(valid.map((s) => ({ id: s.id, priceUsd: Number(s.priceUsd), date: s.saleDate, quantity: s.quantity, isBundle: s.isBundle })))) flagsToApply.push(f); } if (flagsToApply.length) { for (const f of flagsToApply) { await db().update(sales).set({ status: 'flagged', flags: sql`array_append(${sales.flags}, ${f.reason})` }).where(and(eq(sales.id, f.id), eq(sales.status, 'valid'))); const row = saleRows.find((s) => s.id === f.id); if (row) row.status = 'flagged'; } await auditMany(flagsToApply.map((f) => ({ entityType: 'sale', entityId: f.id, action: 'flagged', reason: f.reason, details: { modifiedZ: f.score, assetId } }))); } // 2. observations (guide prices) last 90d per variant const obsRows = await db() .select({ variantId: priceObservations.variantId, priceUsd: priceObservations.priceUsd, date: priceObservations.observationDate, priceKind: priceObservations.priceKind, sourceId: priceObservations.sourceId }) .from(priceObservations) .where(and(eq(priceObservations.assetId, assetId), gte(priceObservations.observationDate, toDateOnly(new Date(now.getTime() - 90 * DAY))))); const obsByVariant = new Map(); for (const o of obsRows) obsByVariant.set(o.variantId ?? '', [...(obsByVariant.get(o.variantId ?? '') ?? []), o]); const premiums = await premiumsFor(asset.categorySlug); const toInput = (s: SaleRow): SaleInput => ({ id: s.id, priceUsd: Number(s.priceUsd), date: s.saleDate, trust: trust.get(s.sourceId) ?? 0.6, confidence: Number(s.confidence), grader: s.grader, grade: s.grade, quantity: s.quantity, isBundle: s.isBundle, status: s.status as SaleInput['status'] }); // 3. per-variant valuations const results: Array<{ variantId: string; isDefault: boolean; out: ReturnType; salesCount: number; sales30d: number; activeL: { count: number; minAsk: number | null } }> = []; const today = toDateOnly(now); // Ask-vs-RIV is recomputed from scratch on every pass: a stale discount computed against a // previous (possibly wrong-variant) valuation must never survive (§84, §196). await db().update(listings).set({ discountToRiv: null }).where(and(eq(listings.assetId, assetId), eq(listings.availability, 'available'), sql`${listings.discountToRiv} is not null`)); for (const v of variants) { const own = (byVariant.get(v.id) ?? []).map(toInput); const comps: NonNullable[0]['comps']> = []; if (own.filter((s) => (s.status ?? 'valid') === 'valid').length < 3 && premiums.length) { for (const other of variants) { if (other.id === v.id) continue; const adj = adjustmentFactor(premiums, { grader: other.grader, grade: other.grade }, { grader: v.grader, grade: v.grade }); if (adj === null) continue; for (const s of byVariant.get(other.id) ?? []) if (s.status === 'valid') comps.push({ priceUsd: Number(s.priceUsd), date: s.saleDate, adjustment: adj, trust: trust.get(s.sourceId) }); } } const out = computeValuation({ sales: own, comps, observations: (obsByVariant.get(v.id) ?? []).map((o) => ({ priceUsd: Number(o.priceUsd), date: new Date(`${o.date}T00:00:00Z`), priceKind: o.priceKind, trust: trust.get(o.sourceId) })), now, categorySlug: asset.categorySlug, }); const salesCount = own.filter((s) => (s.status ?? 'valid') === 'valid').length; const sales30d = own.filter((s) => (s.status ?? 'valid') === 'valid' && now.getTime() - s.date.getTime() <= 30 * DAY).length; const activeL = await activeListingStats(assetId, v.id); // §35: say when the inputs include an estimated (not invoiced) buyer premium if (out.salesUsed.length) { const used = new Set(out.salesUsed); const est = (byVariant.get(v.id) ?? []).filter((s) => used.has(s.id) && s.feeBasis?.startsWith('added_')).length; if (est) out.notes.push(`${Math.round((100 * est) / out.salesUsed.length)} % of inputs include an estimated buyer premium`); } results.push({ variantId: v.id, isDefault: v.isDefault, out, salesCount, sales30d, activeL }); if (out.riv === null && salesCount === 0 && (obsByVariant.get(v.id) ?? []).length === 0) continue; await db().insert(valuations).values({ id: newId('valuation'), assetId, variantId: v.id, computedAt: now, rivUsd: out.riv, lowUsd: out.low, highUsd: out.high, confidence: out.confidence, confidenceLabel: out.label, sampleSize: out.sampleSize, windowDays: out.windowDays, methods: out.methods, salesUsed: out.salesUsed, observationsUsed: out.observationsUsed, method: `ensemble_v1:${out.basis}`, notes: out.notes, }); const prev30 = await snapshotValue(assetId, v.id, new Date(now.getTime() - 30 * DAY)); const prev1y = await snapshotValue(assetId, v.id, new Date(now.getTime() - 365 * DAY)); const latest = own[0]; await db() .insert(variantStats) .values({ variantId: v.id, assetId, rivUsd: out.riv, rivLowUsd: out.low, rivHighUsd: out.high, rivConfidence: out.confidence, rivSampleSize: out.sampleSize, latestSaleUsd: latest?.priceUsd ?? null, latestSaleAt: latest?.date ?? null, change30d: safeRatioV(pctChange(prev30, out.riv)), change1y: safeRatioV(pctChange(prev1y, out.riv)), salesCount, sales30d, activeListings: activeL.count, minAskUsd: activeL.minAsk, liquidityScore: liquidityScore({ salesPerMonth: (own.length / 36) || 0, activeListings: activeL.count, sources: new Set((byVariant.get(v.id) ?? []).map((s) => s.sourceId)).size, medianDaysBetweenSales: medianGapDays(own.map((s) => s.date)), askSoldSpread: activeL.minAsk && out.riv ? (activeL.minAsk - out.riv) / out.riv : null }), updatedAt: now }) .onConflictDoUpdate({ target: variantStats.variantId, set: { rivUsd: out.riv, rivLowUsd: out.low, rivHighUsd: out.high, rivConfidence: out.confidence, rivSampleSize: out.sampleSize, latestSaleUsd: latest?.priceUsd ?? null, latestSaleAt: latest?.date ?? null, change30d: safeRatioV(pctChange(prev30, out.riv)), change1y: safeRatioV(pctChange(prev1y, out.riv)), salesCount, sales30d, activeListings: activeL.count, minAskUsd: activeL.minAsk, updatedAt: now } }); await db() .insert(priceSnapshots) .values({ assetId, variantId: v.id, date: today, rivUsd: out.riv, latestSaleUsd: latest?.priceUsd ?? null, medianUsd: out.distribution.median, salesCount: own.filter((s) => toDateOnly(s.date) === today).length, volumeUsd: own.filter((s) => toDateOnly(s.date) === today).reduce((a, s) => a + s.priceUsd, 0) || null, listingsCount: activeL.count, minAskUsd: activeL.minAsk, observationUsd: out.methods.guide }) .onConflictDoUpdate({ target: [priceSnapshots.assetId, priceSnapshots.variantId, priceSnapshots.date], set: { rivUsd: out.riv, latestSaleUsd: latest?.priceUsd ?? null, medianUsd: out.distribution.median, listingsCount: activeL.count, minAskUsd: activeL.minAsk, observationUsd: out.methods.guide } }); // Ask vs RIV on the active listings of THIS variant only (§84 data-quality gate): // – the valuation must rest on ≥ ASK_MIN_SAMPLE transactions with confidence ≥ ASK_MIN_CONFIDENCE // (comps-only / guide-only estimates never qualify an ask); // – the listing must be confidently matched and carry a readable grade when it is slabbed; // – an ask outside [0.1×, 10×] RIV is flagged `riv_anomaly` (identity/data problem), not priced. if (out.riv !== null && out.basis === 'transactions' && out.confidence >= ASK_MIN_CONFIDENCE && out.sampleSize >= ASK_MIN_SAMPLE) { const lo = out.riv * ASK_ANOMALY_LOW_RATIO; const hi = out.riv * ASK_ANOMALY_HIGH_RATIO; // Auction lots are excluded: a current/opening bid (often ¥1 on Yahoo! Auctions) is not an asking price. // Bid-vs-RIV belongs to auction intelligence with fees (§33–§35), not to the deal rails. const scope = and(eq(listings.variantId, v.id), eq(listings.availability, 'available'), sql`${listings.priceUsd} > 0`, sql`${listings.listingType} <> 'auction'`); await db() .update(listings) .set({ discountToRiv: sql`round((${listings.priceUsd} - ${out.riv}) / ${out.riv}, 4)`, // > 50 % below RIV (§174): keep the number, flag `riv_review`, never surface as a deal flags: sql`case when ${listings.priceUsd} < ${out.riv * (1 + DEAL_REVIEW_THRESHOLD)} then (case when 'riv_review' = any(${listings.flags}) then array_remove(${listings.flags}, 'riv_anomaly') else array_append(array_remove(${listings.flags}, 'riv_anomaly'), 'riv_review') end) else array_remove(array_remove(${listings.flags}, 'riv_anomaly'), 'riv_review') end`, }) .where(and(scope, sql`${listings.priceUsd} between ${lo} and ${hi}`, sql`${listings.confidence} >= ${ASK_MIN_MATCH_CONFIDENCE}`, sql`not (${listings.grader} is not null and ${listings.grade} is null)`)); await db() .update(listings) .set({ discountToRiv: null, flags: sql`case when 'riv_anomaly' = any(${listings.flags}) then ${listings.flags} else array_append(${listings.flags}, 'riv_anomaly') end` }) .where(and(scope, sql`(${listings.priceUsd} < ${lo} or ${listings.priceUsd} > ${hi})`)); } } // 4. asset-level representative variant (the headline RIV). Preference order: // transaction-based valuation with confidence ≥ 0.5 and ≥ 5 sales → the DEFAULT (raw/base) variant // when it qualifies, else the qualifying variant with the most sales; fallback: any priced variant // by confidence. Rationale: a graded sub-variant with a tighter distribution used to win over the // base market and made the headline value (and every ask comparison) wrong by an order of magnitude. const priced = results.filter((r) => r.out.riv !== null).sort((a, b) => b.out.confidence - a.out.confidence || b.salesCount - a.salesCount); const qualified = priced.filter((r) => r.out.basis === 'transactions' && r.out.confidence >= ASK_MIN_CONFIDENCE && r.out.sampleSize >= ASK_MIN_SAMPLE); const rep = qualified.find((r) => r.isDefault) ?? qualified.sort((a, b) => b.salesCount - a.salesCount || b.out.confidence - a.out.confidence)[0] ?? priced[0] ?? null; const validSales = saleRows.filter((s) => s.status === 'valid'); const latest = validSales[0]; // Changes are measured on the representative variant's own series so that a change of // representative variant never shows up as a price move. const repSeries = rep?.variantId ?? ''; const prevAsset = { d1: await snapshotValue(assetId, repSeries, new Date(now.getTime() - 1 * DAY)), d7: await snapshotValue(assetId, repSeries, new Date(now.getTime() - 7 * DAY)), d30: await snapshotValue(assetId, repSeries, new Date(now.getTime() - 30 * DAY)), d90: await snapshotValue(assetId, repSeries, new Date(now.getTime() - 90 * DAY)), y1: await snapshotValue(assetId, repSeries, new Date(now.getTime() - 365 * DAY)) }; const riv = rep?.out.riv ?? null; const activeAll = await activeListingStats(assetId, null); const repActive = rep?.activeL ?? null; // ATH / ATL / drawdown belong to the same series as the headline RIV: the representative variant. // (An asset-wide ATH mixed a sealed copy's record with a loose copy's valuation → −99 % "drawdowns".) const extremes = rep ? (await db().select({ n: sql`count(*)::int`, min: sql`min(${BUYER_PAYS})`, max: sql`max(${BUYER_PAYS})`, minAt: sql`(array_agg(${sales.saleDate} order by ${BUYER_PAYS} asc))[1]`, maxAt: sql`(array_agg(${sales.saleDate} order by ${BUYER_PAYS} desc))[1]` }).from(sales).where(and(eq(sales.assetId, assetId), eq(sales.variantId, rep.variantId), eq(sales.status, 'valid'), sql`${sales.quantity} = 1`, sql`not ${sales.isBundle}`)))[0]! : allSalesCount; const salesCount = allSalesCount.n; const sales30d = validSales.filter((s) => now.getTime() - s.saleDate.getTime() <= 30 * DAY).length; const sales1y = validSales.filter((s) => now.getTime() - s.saleDate.getTime() <= 365 * DAY).length; const salesPrev30 = validSales.filter((s) => { const age = now.getTime() - s.saleDate.getTime(); return age > 30 * DAY && age <= 60 * DAY; }).length; const pop = await latestPopulation(assetId); const [obsCount] = await db().select({ n: sql`count(*)::int` }).from(priceObservations).where(eq(priceObservations.assetId, assetId)); const sourcesCount = new Set([...saleRows.map((s) => s.sourceId), ...obsRows.map((o) => o.sourceId)]).size; const liq = liquidityScore({ salesPerMonth: sales1y / 12, activeListings: activeAll.count, sources: sourcesCount, medianDaysBetweenSales: medianGapDays(validSales.map((s) => s.saleDate)), askSoldSpread: repActive?.minAsk && riv ? (repActive.minAsk - riv) / riv : null }); const rar = rarityScore({ population: pop?.total ?? null, productionQuantity: asset.productionQuantity, listingsPerYear: activeAll.count > 0 || sales1y > 0 ? activeAll.count * 4 : null, salesPerYear: sales1y > 0 || salesCount > 0 ? sales1y : null, populationGrowthPct: pop?.growth ?? null }); // ratio columns are numeric(8,6): anything ≥ ±99 (9,900 %) would overflow and abort the asset; such a // move is never a market signal but a data/identity artefact → stored as null (§143, §196). const safeRatio = (v: number | null): number | null => (v === null || !Number.isFinite(v) || Math.abs(v) >= 99 ? null : v); const ch = { d1: safeRatio(pctChange(prevAsset.d1, riv)), d7: safeRatio(pctChange(prevAsset.d7, riv)), d30: safeRatio(pctChange(prevAsset.d30, riv)), d90: safeRatio(pctChange(prevAsset.d90, riv)), y1: safeRatio(pctChange(prevAsset.y1, riv)) }; const mom = { m7: momentumScore({ priceChange: ch.d7, volumeNow: validSales.filter((s) => now.getTime() - s.saleDate.getTime() <= 7 * DAY).length, volumePrev: validSales.filter((s) => { const a = now.getTime() - s.saleDate.getTime(); return a > 7 * DAY && a <= 14 * DAY; }).length }), m30: momentumScore({ priceChange: ch.d30, volumeNow: sales30d, volumePrev: salesPrev30 }), m90: momentumScore({ priceChange: ch.d90, volumeNow: validSales.filter((s) => now.getTime() - s.saleDate.getTime() <= 90 * DAY).length, volumePrev: validSales.filter((s) => { const a = now.getTime() - s.saleDate.getTime(); return a > 90 * DAY && a <= 180 * DAY; }).length }), y1: momentumScore({ priceChange: ch.y1, volumeNow: sales1y, volumePrev: validSales.length - sales1y }) }; const listingMomentum = activeAll.count > 0 ? Math.min(100, activeAll.count * 5) : null; const trending = trendingScore({ priceMomentum: mom.m30, volumeMomentum: momentumScore({ priceChange: null, volumeNow: sales30d, volumePrev: salesPrev30 }), searchMomentum: null, listingMomentum, newsMomentum: null }); // Best gated ask across the asset's variants (negative = below its own variant's RIV). No fallback: // comparing the cheapest ask of ANY variant against the headline RIV was the −8,500 % bug. const [bestListing] = await db().select({ d: sql`min(${listings.discountToRiv})` }).from(listings).where(and(eq(listings.assetId, assetId), eq(listings.availability, 'available'), sql`${listings.discountToRiv} is not null`, sql`not ('riv_review' = any(${listings.flags}))`)); const fields = { brand: asset.brand, set: asset.setName, number: asset.number, year: asset.year, variant: asset.variant, image: asset.heroImageUrl, description: asset.description, identifiers: Object.keys(asset.identifiers).length ? 1 : null }; const dq = assetDataQuality({ fieldsPresent: Object.values(fields).filter((x) => x !== null && x !== undefined).length, fieldsTotal: Object.keys(fields).length, sourceTrustAvg: sourcesCount ? median([...new Set(saleRows.map((s) => s.sourceId))].map((s) => trust.get(s) ?? 0.5)) : null, identificationConfidence: rep ? Number(median(validSales.map((s) => Number(s.confidence))) ?? 0.8) : null, hasImage: Boolean(asset.heroImageUrl), salesCount }); const statsRow = { assetId, rivUsd: riv, rivLowUsd: rep?.out.low ?? null, rivHighUsd: rep?.out.high ?? null, rivConfidence: rep?.out.confidence ?? null, rivSampleSize: rep?.out.sampleSize ?? 0, rivVariantId: rep?.variantId ?? null, latestSaleUsd: latest ? Number(latest.priceUsd) : null, latestSaleAt: latest?.saleDate ?? null, change1d: ch.d1, change7d: ch.d7, change30d: ch.d30, change90d: ch.d90, change1y: ch.y1, athUsd: extremes.n ? Number(extremes.max) : null, athAt: extremes.n ? asDate(extremes.maxAt) : null, atlUsd: extremes.n ? Number(extremes.min) : null, atlAt: extremes.n ? asDate(extremes.minAt) : null, salesCount, sales30d, sales1y, volume30dUsd: sales30d ? validSales.filter((s) => now.getTime() - s.saleDate.getTime() <= 30 * DAY).reduce((a, s) => a + Number(s.priceUsd), 0) : null, activeListings: activeAll.count, // lowest ask of the representative variant (comparable to the headline RIV); all variants otherwise minAskUsd: repActive?.minAsk ?? (rep ? null : activeAll.minAsk), // never a graded slab's ask next to a raw RIV observationsCount: obsCount?.n ?? 0, sourcesCount, liquidityScore: liq, rarityScore: rar, momentum7d: mom.m7, momentum30d: mom.m30, momentum90d: mom.m90, momentum1y: mom.y1, trendingScore: trending, valueOpportunity: bestListing?.d !== null && bestListing?.d !== undefined ? Number(bestListing.d) : null, dataQuality: dq, updatedAt: now, }; await db().insert(assetStats).values(statsRow).onConflictDoUpdate({ target: assetStats.assetId, set: { ...statsRow, watchers: sql`${assetStats.watchers}`, views30d: sql`${assetStats.views30d}` } }); await db().update(assets).set({ dataQuality: dq, updatedAt: now }).where(eq(assets.id, assetId)); await db() .insert(priceSnapshots) .values({ assetId, variantId: '', date: today, rivUsd: riv, latestSaleUsd: latest ? Number(latest.priceUsd) : null, medianUsd: rep?.out.distribution.median ?? null, salesCount: validSales.filter((s) => toDateOnly(s.saleDate) === today).length, volumeUsd: null, listingsCount: activeAll.count, minAskUsd: activeAll.minAsk, observationUsd: rep?.out.methods.guide ?? null }) .onConflictDoUpdate({ target: [priceSnapshots.assetId, priceSnapshots.variantId, priceSnapshots.date], set: { rivUsd: riv, latestSaleUsd: latest ? Number(latest.priceUsd) : null, listingsCount: activeAll.count, minAskUsd: activeAll.minAsk } }); if (opts.rebuildHistory) await rebuildSnapshotHistory(assetId, saleRows, variants.map((v) => v.id), now); if (riv !== null) await emit('valuation_updated', { type: 'asset', id: assetId }, { riv, confidence: rep?.out.confidence, sampleSize: rep?.out.sampleSize }); return { variants: results.length, riv }; } async function snapshotValue(assetId: string, variantId: string, at: Date): Promise { const [row] = await db() .select({ v: priceSnapshots.rivUsd }) .from(priceSnapshots) .where(and(eq(priceSnapshots.assetId, assetId), eq(priceSnapshots.variantId, variantId), sql`${priceSnapshots.date} <= ${toDateOnly(at)}`, sql`${priceSnapshots.rivUsd} is not null`)) .orderBy(desc(priceSnapshots.date)) .limit(1); return row?.v === null || row?.v === undefined ? null : Number(row.v); } async function activeListingStats(assetId: string, variantId: string | null): Promise<{ count: number; minAsk: number | null }> { const [row] = await db() // count = every live listing; min ask = fixed-price / best-offer / ask only (an auction bid is not an ask) .select({ n: sql`count(*)::int`, min: sql`min(${listings.priceUsd}) filter (where ${listings.listingType} <> 'auction' and ${listings.priceUsd} > 0)` }) .from(listings) .where(and(eq(listings.assetId, assetId), eq(listings.availability, 'available'), variantId ? eq(listings.variantId, variantId) : undefined)); return { count: row?.n ?? 0, minAsk: row?.min === null || row?.min === undefined ? null : Number(row.min) }; } async function latestPopulation(assetId: string): Promise<{ total: number; growth: number | null } | null> { const rows = await db().select({ total: populationReports.total, date: populationReports.reportDate }).from(populationReports).where(eq(populationReports.assetId, assetId)).orderBy(desc(populationReports.reportDate)).limit(2); if (!rows.length) return null; const growth = rows.length === 2 && rows[1]!.total > 0 ? (rows[0]!.total - rows[1]!.total) / rows[1]!.total : null; return { total: rows[0]!.total, growth }; } function medianGapDays(dates: Date[]): number | null { if (dates.length < 2) return null; const s = [...dates].sort((a, b) => a.getTime() - b.getTime()); const gaps: number[] = []; for (let i = 1; i < s.length; i++) gaps.push((s[i]!.getTime() - s[i - 1]!.getTime()) / DAY); return median(gaps); } /** * Historical snapshots (weekly) reconstructed from transactions only, computing the valuation * as of each week with the sales known up to that date. Gives asset charts and index history a * past without inventing observations. */ export async function rebuildSnapshotHistory(assetId: string, saleRows: SaleRow[], variantIds: string[], now: Date): Promise { const valid = saleRows.filter((s) => s.status === 'valid').sort((a, b) => a.saleDate.getTime() - b.saleDate.getTime()); if (valid.length < 2) return 0; const start = valid[0]!.saleDate.getTime(); const rows: Array = []; const trust = await trustMap(); for (let t = start + 7 * DAY; t < now.getTime() - DAY; t += 7 * DAY) { const asOf = new Date(t); const date = toDateOnly(asOf); const upTo = valid.filter((s) => s.saleDate.getTime() <= t); // asset-level: representative = variant with most sales up to date const counts = new Map(); for (const s of upTo) counts.set(s.variantId ?? '', (counts.get(s.variantId ?? '') ?? 0) + 1); for (const vid of [...variantIds, '']) { const subset = vid === '' ? upTo.filter((s) => (s.variantId ?? '') === [...counts.entries()].sort((a, b) => b[1] - a[1])[0]?.[0]) : upTo.filter((s) => s.variantId === vid); if (subset.length < 2) continue; const out = computeValuation({ sales: subset.map((s) => ({ id: s.id, priceUsd: Number(s.priceUsd), date: s.saleDate, trust: trust.get(s.sourceId) ?? 0.6, confidence: Number(s.confidence) })), now: asOf }); if (out.riv === null) continue; const weekSales = subset.filter((s) => t - s.saleDate.getTime() < 7 * DAY); rows.push({ assetId, variantId: vid, date, rivUsd: out.riv, latestSaleUsd: subset.at(-1)!.priceUsd, medianUsd: out.distribution.median, salesCount: weekSales.length, volumeUsd: weekSales.reduce((a, s) => a + Number(s.priceUsd), 0) || null, listingsCount: 0, minAskUsd: null, observationUsd: null }); } } for (let i = 0; i < rows.length; i += 500) { await db() .insert(priceSnapshots) .values(rows.slice(i, i + 500)) .onConflictDoUpdate({ target: [priceSnapshots.assetId, priceSnapshots.variantId, priceSnapshots.date], set: { rivUsd: sql`excluded.riv_usd`, medianUsd: sql`excluded.median_usd`, latestSaleUsd: sql`excluded.latest_sale_usd`, salesCount: sql`excluded.sales_count`, volumeUsd: sql`excluded.volume_usd` } }); } return rows.length; } /** Assets whose evidence changed since their last stats update (or all). */ export async function assetsNeedingValuation(opts: { all?: boolean; limit?: number } = {}): Promise { const limit = opts.limit ?? 5000; if (opts.all) { const rows = await db().select({ id: assets.id }).from(assets).orderBy(assets.createdAt).limit(limit); return rows.map((r) => r.id); } const rows = await db().execute(sql` with touched as ( select asset_id, max(created_at) as t from sales group by asset_id union all select asset_id, max(updated_at) from listings group by asset_id union all select asset_id, max(created_at) from price_observations group by asset_id ), agg as (select asset_id, max(t) as t from touched group by asset_id) select a.asset_id from agg a left join asset_stats s on s.asset_id = a.asset_id where s.asset_id is null or s.updated_at < a.t limit ${limit}`); return (rows as unknown as Array<{ asset_id: string }>).map((r) => r.asset_id); } export async function valueMany(ids: string[], opts: { rebuildHistory?: boolean; concurrency?: number } = {}): Promise<{ valued: number; priced: number }> { let valued = 0; let priced = 0; const conc = opts.concurrency ?? 4; for (let i = 0; i < ids.length; i += conc) { const chunk = ids.slice(i, i + conc); const res = await Promise.all(chunk.map((id) => valueAsset(id, { rebuildHistory: opts.rebuildHistory }).catch((err) => { log.error({ err, id }, 'valueAsset failed'); return null; }))); for (const r of res) { if (!r) continue; valued++; if (r.riv !== null) priced++; } if ((i / conc) % 50 === 0 && i > 0) log.info({ done: i, total: ids.length, priced }, 'valuation progress'); } return { valued, priced }; } /** Empirical grade premiums per category from sales in the last 2 years (§118). */ export async function computePremiums(categorySlug?: string): Promise { const since = new Date(Date.now() - 2 * 365 * DAY); const cats = categorySlug ? [categorySlug] : (await db().selectDistinct({ c: assets.categorySlug }).from(assets)).map((r) => r.c); let written = 0; for (const cat of cats) { const rows = await db() .select({ assetId: sales.assetId, grader: sales.grader, grade: sales.grade, priceUsd: BUYER_PAYS }) .from(sales) .innerJoin(assets, eq(assets.id, sales.assetId)) .where(and(eq(assets.categorySlug, cat), eq(sales.status, 'valid'), gte(sales.saleDate, since))); const prem = computeGradePremiums(rows.map((r) => ({ assetId: r.assetId, grader: r.grader, grade: r.grade, priceUsd: Number(r.priceUsd) }))); for (const p of prem) { await db() .insert(gradePremiums) .values({ id: newId('valuation'), categorySlug: cat, grader: p.grader, grade: p.grade, marketMultiplier: p.marketMultiplier, sampleSize: p.sampleSize, computedAt: new Date() }) .onConflictDoUpdate({ target: [gradePremiums.categorySlug, gradePremiums.grader, gradePremiums.grade], set: { marketMultiplier: p.marketMultiplier, sampleSize: p.sampleSize, computedAt: new Date() } }); written++; } premiumCache.delete(cat); } return written; } export { inArray, isNull, or };