TypeScript 61.9%
HTML 37.2%
SQL 0.7%
1import { and, eq, sql } from 'drizzle-orm';2import { auctionLots, auctions, images, listingEvents, listings, news, populationReports, priceObservations, sales, sources } from '@rareindex/database';3import { normalizeCondition, parseGradeFromTitle } from '@rareindex/taxonomy';4import { allInPrice, feeScheduleFor } from '@rareindex/valuation';5import type { Grade } from '@rareindex/shared';6import { newId, sha256, toDateOnly, type NormalizedAuctionLot, type NormalizedCatalogItem, type NormalizedListing, type NormalizedNewsItem, type NormalizedPopulationReport, type NormalizedPriceObservation, type NormalizedRecord, type NormalizedSale } from '@rareindex/shared';7import { db } from '../lib/db.ts';8import { toUsd, usdRateFor } from '../lib/fx.ts';9import { audit } from '../lib/audit.ts';10import { enrichAttributesFromTitle, ensureVariant, resolveAsset, type Resolution } from './resolver.ts';11import { recordCertificate } from './certificates.ts';1213/** Record a slab certification sighting when a record carries one (SPEC §22); never throws into the writer. */14async function trackCert(rec: { grade: { grader: string | null; grade: string | null; qualifier: string | null; certificationNumber: string | null }; sourceId: string; connectorId: string; sourceUrl: string; rawTitle?: string | null }, kind: 'sale' | 'listing' | 'auction_lot', targetId: string | null, assetId: string | null, variantId: string | null, priceUsd: number | null, price: number | null, currency: string | null, observedAt: Date): Promise<void> {15 const g = refineGrade(rec.grade, rec.rawTitle);16 if (!g.certificationNumber || !g.grader) return;17 try {18 await recordCertificate({ grader: g.grader, certNumber: g.certificationNumber, assetId, variantId, grade: g.grade, qualifier: g.qualifier, kind, targetId, sourceId: rec.sourceId, connectorId: rec.connectorId, sourceUrl: rec.sourceUrl, priceUsd, price, currency, observedAt });19 } catch {20 /* cert tracking is best-effort; the canonical write already succeeded */21 }22}2324export class FxMissingError extends Error {25 constructor(public currency: string, public date: Date) {26 super(`fx_missing:${currency}:${toDateOnly(date)}`);27 }28}2930export interface ApplyResult {31 assetId: string | null;32 variantId: string | null;33 targetId: string | null;34 method: Resolution['method'] | null;35 confidence: number | null;36 status: 'applied' | 'unmatched' | 'rejected' | 'duplicate';37 reason?: string;38 event?: 'sale_detected' | 'listing_created' | 'listing_updated' | 'entity_created';39}4041const sourceMetaCache = new Map<string, { name: string; sourceType: string }>();42async function sourceMeta(sourceId: string): Promise<{ name: string; sourceType: string } | null> {43 const hit = sourceMetaCache.get(sourceId);44 if (hit) return hit;45 const [row] = await db().select({ name: sources.name, sourceType: sources.sourceType }).from(sources).where(eq(sources.id, sourceId)).limit(1);46 if (!row) return null;47 sourceMetaCache.set(sourceId, row);48 return row;49}5051/**52 * Buyer-pays price of a sale (§35): price_usd + the house's buyer premium when the record is53 * hammer-only. `fxToScheduleCurrency` converts the native amount into the schedule's currency so54 * marginal tiers apply at the right thresholds. Returns nulls when the basis cannot be established.55 */56export async function saleAllIn(rec: { priceUsd: number; currency: string; buyerPremiumIncluded: boolean | null | undefined; auctionHouse: string | null | undefined; saleType: string | null | undefined; sourceId: string; connectorId: string; saleDate: Date }): Promise<{ allInUsd: number; feeBasis: string; buyerPremiumRate: number }> {57 const meta = await sourceMeta(rec.sourceId);58 const house = rec.auctionHouse ?? meta?.name ?? rec.connectorId;59 const schedule = feeScheduleFor(house) ?? feeScheduleFor(rec.connectorId);60 let fxToScheduleCurrency = 1;61 if (schedule && schedule.currency !== 'USD') {62 // price is already USD here; tiers are in the schedule currency → USD × (schedule units per USD)63 const r = await usdRateFor(schedule.currency, rec.saleDate);64 if (r && r.rate > 0) fxToScheduleCurrency = r.rate;65 }66 const a = allInPrice({ price: rec.priceUsd, buyerPremiumIncluded: rec.buyerPremiumIncluded, house: schedule ? schedule.id : house, saleType: rec.saleType, sourceType: meta?.sourceType ?? null, fxToScheduleCurrency });67 return { allInUsd: a.allIn, feeBasis: a.basis, buyerPremiumRate: a.rate };68}6970const trustCache = new Map<string, number>();71async function sourceTrust(sourceId: string): Promise<number> {72 const hit = trustCache.get(sourceId);73 if (hit !== undefined) return hit;74 const [row] = await db().select({ t: sources.trustScore }).from(sources).where(eq(sources.id, sourceId)).limit(1);75 const t = row ? Number(row.t) : 0.5;76 trustCache.set(sourceId, t);77 return t;78}7980/** Data quality 0–100 (§150): completeness × source trust × identification confidence. */81function dataQuality(fields: Record<string, unknown>, trust: number, confidence: number, extra = 0): number {82 const keys = Object.keys(fields);83 const present = keys.filter((k) => fields[k] !== null && fields[k] !== undefined && fields[k] !== '' && !(Array.isArray(fields[k]) && (fields[k] as unknown[]).length === 0)).length;84 const completeness = keys.length ? present / keys.length : 0;85 return Math.round(Math.max(0, Math.min(100, (0.4 * completeness + 0.3 * trust + 0.3 * confidence + extra) * 100)));86}8788function normCond(categorySlug: string, raw: string | null | undefined, already: string | null | undefined): string | null {89 return already ?? normalizeCondition(categorySlug, raw);90}9192async function storeImages(assetId: string, urls: string[], sourceId: string, ref: { listingId?: string; saleId?: string }): Promise<void> {93 if (!urls.length) return;94 const rows = urls.slice(0, 12).map((url, i) => ({ id: newId('image'), assetId, listingId: ref.listingId ?? null, saleId: ref.saleId ?? null, sourceId, url, role: i === 0 && !ref.listingId && !ref.saleId ? 'hero' : 'source' }));95 await db().insert(images).values(rows).onConflictDoNothing({ target: images.url });96}9798export async function applyRecord(rec: NormalizedRecord): Promise<ApplyResult> {99 switch (rec.kind) {100 case 'catalog_item':101 return applyCatalog(rec);102 case 'sale':103 return applySale(rec);104 case 'listing':105 return applyListing(rec);106 case 'price_observation':107 return applyObservation(rec);108 case 'auction_lot':109 return applyAuctionLot(rec);110 case 'population_report':111 return applyPopulation(rec);112 case 'news_item':113 return applyNews(rec);114 }115}116117/**118 * Grade refinement (§18, §84): when a connector did not read the slab (or read the grader but not the119 * grade), parse the raw title. A grader with no readable grade stays "grader · grade unknown" — it must120 * never fall into the raw/ungraded variant, which is what compared PSA slabs against raw valuations.121 * An explicit `raw` from the connector is trusted over the title.122 */123export function refineGrade(grade: Grade, rawTitle: string | null | undefined): Grade {124 if (!rawTitle || grade.grader === 'raw' || (grade.grader && grade.grade)) return grade;125 const parsed = parseGradeFromTitle(rawTitle);126 if (!parsed.grader || parsed.grader === 'raw') return grade;127 if (grade.grader && parsed.grader !== grade.grader) return grade; // never contradict the connector's grader128 // A grade recorded without a grader is a condition-scale value (e.g. a price-guide tab), not a slab129 // grade: when the title reveals the slab, the title's grade wins.130 const g = grade.grader ? (grade.grade ?? parsed.grade) : parsed.grade;131 return { ...grade, grader: grade.grader ?? parsed.grader, grade: g, qualifier: grade.qualifier ?? parsed.qualifier };132}133134async function applyCatalog(rec: NormalizedCatalogItem): Promise<ApplyResult> {135 const r = await resolveAsset(rec.attributes, { enrich: true, imageUrls: rec.imageUrls, description: rec.description });136 if (!r) return { assetId: null, variantId: null, targetId: null, method: null, confidence: null, status: 'unmatched' };137 if (rec.imageUrls.length) await storeImages(r.assetId, rec.imageUrls.slice(0, 1), rec.sourceId, {});138 return { assetId: r.assetId, variantId: null, targetId: r.assetId, method: r.method, confidence: r.confidence, status: 'applied', event: r.created ? 'entity_created' : undefined };139}140141async function applySale(rec: NormalizedSale): Promise<ApplyResult> {142 const grade = refineGrade(rec.grade, rec.rawTitle);143 const r = await resolveAsset(enrichAttributesFromTitle(rec.attributes, rec.rawTitle), { imageUrls: rec.imageUrls });144 if (!r) return { assetId: null, variantId: null, targetId: null, method: null, confidence: null, status: 'unmatched' };145 const condition = normCond(rec.attributes.categorySlug, rec.condition.conditionRaw, rec.condition.condition);146 const variant = await ensureVariant(r.assetId, grade, { condition, completeness: rec.condition.completeness }, rec.attributes.size);147 const fx = await toUsd(rec.price, rec.currency, rec.saleDate);148 if (!fx) throw new FxMissingError(rec.currency, rec.saleDate);149 const dedupeKey = sha256(['sale', rec.sourceId, rec.externalId ?? rec.sourceUrl, toDateOnly(rec.saleDate), rec.price.toFixed(2)].join('|'));150 const trust = await sourceTrust(rec.sourceId);151 const fees = await saleAllIn({ priceUsd: fx.usd, currency: rec.currency, buyerPremiumIncluded: rec.buyerPremiumIncluded, auctionHouse: rec.auctionHouse, saleType: rec.saleType, sourceId: rec.sourceId, connectorId: rec.connectorId, saleDate: rec.saleDate });152 const confidence = Math.min(rec.confidence, r.confidence);153 const flags: string[] = [];154 if (rec.isBundle || rec.quantity > 1) flags.push('bundle');155 if (rec.price <= 0) flags.push('zero_price');156 if (confidence < 0.6) flags.push('low_identification_confidence');157 const status = flags.includes('bundle') || flags.includes('zero_price') ? 'excluded' : 'valid';158 const id = newId('sale');159 const inserted = await db()160 .insert(sales)161 .values({162 id,163 assetId: r.assetId,164 variantId: variant.id,165 sourceId: rec.sourceId,166 connectorId: rec.connectorId,167 sourceUrl: rec.sourceUrl,168 externalId: rec.externalId,169 saleType: rec.saleType,170 saleDate: rec.saleDate,171 price: rec.price,172 currency: rec.currency,173 priceUsd: fx.usd,174 fxRate: fx.rate,175 fxDate: fx.fxDate,176 buyerPremiumIncluded: rec.buyerPremiumIncluded,177 allInUsd: fees.allInUsd,178 feeBasis: fees.feeBasis,179 buyerPremiumRate: fees.buyerPremiumRate,180 quantity: rec.quantity,181 isBundle: rec.isBundle,182 condition,183 grader: grade.grader,184 grade: grade.grade,185 certificationNumber: grade.certificationNumber,186 location: rec.location,187 auctionHouse: rec.auctionHouse,188 lotNumber: rec.lotNumber,189 imageUrls: rec.imageUrls,190 rawTitle: rec.rawTitle,191 confidence,192 dataQuality: dataQuality({ date: rec.saleDate, price: rec.price, currency: rec.currency, images: rec.imageUrls, grade: grade.grade, condition, externalId: rec.externalId, type: rec.saleType === 'unknown' ? null : rec.saleType }, trust, confidence),193 status,194 flags,195 dedupeKey,196 })197 .onConflictDoNothing({ target: sales.dedupeKey })198 .returning({ id: sales.id });199 if (inserted.length === 0) return { assetId: r.assetId, variantId: variant.id, targetId: null, method: r.method, confidence, status: 'duplicate' };200 if (status === 'excluded') await audit({ entityType: 'sale', entityId: id, action: 'excluded', reason: flags.join(','), details: { rawTitle: rec.rawTitle, quantity: rec.quantity } });201 if (rec.imageUrls.length) await storeImages(r.assetId, rec.imageUrls.slice(0, 3), rec.sourceId, { saleId: id });202 await trackCert(rec, 'sale', id, r.assetId, variant.id, fx.usd, rec.price, rec.currency, rec.saleDate);203 return { assetId: r.assetId, variantId: variant.id, targetId: id, method: r.method, confidence, status: 'applied', event: 'sale_detected' };204}205206async function applyListing(rec: NormalizedListing): Promise<ApplyResult> {207 const grade = refineGrade(rec.grade, rec.rawTitle);208 const r = await resolveAsset(enrichAttributesFromTitle(rec.attributes, rec.rawTitle), { imageUrls: rec.imageUrls });209 if (!r) return { assetId: null, variantId: null, targetId: null, method: null, confidence: null, status: 'unmatched' };210 const condition = normCond(rec.attributes.categorySlug, rec.condition.conditionRaw, rec.condition.condition);211 const variant = await ensureVariant(r.assetId, grade, { condition, completeness: rec.condition.completeness }, rec.attributes.size);212 let priceUsd: number | null = null;213 if (rec.price !== null && rec.currency) {214 const fx = await toUsd(rec.price, rec.currency, rec.observedAt);215 if (!fx) throw new FxMissingError(rec.currency, rec.observedAt);216 priceUsd = fx.usd;217 }218 const externalId = rec.externalId ?? sha256(rec.sourceUrl);219 const trust = await sourceTrust(rec.sourceId);220 const confidence = Math.min(rec.confidence, r.confidence);221 const now = rec.observedAt;222 const [existing] = await db().select({ id: listings.id, price: listings.price, availability: listings.availability }).from(listings).where(and(eq(listings.sourceId, rec.sourceId), eq(listings.externalId, externalId))).limit(1);223 const quality = dataQuality({ price: rec.price, currency: rec.currency, images: rec.imageUrls, seller: rec.seller, condition, grade: grade.grade, listedAt: rec.listedAt, description: rec.description }, trust, confidence);224 if (!existing) {225 const id = newId('listing');226 await db().insert(listings).values({227 id,228 assetId: r.assetId,229 variantId: variant.id,230 sourceId: rec.sourceId,231 connectorId: rec.connectorId,232 sourceUrl: rec.sourceUrl,233 externalId,234 listingType: rec.listingType,235 price: rec.price,236 currency: rec.currency,237 priceUsd,238 seller: rec.seller,239 sellerReputation: rec.sellerReputation,240 location: rec.location,241 shippingCost: rec.shippingCost,242 quantity: rec.quantity,243 condition,244 grader: grade.grader,245 grade: grade.grade,246 certificationNumber: grade.certificationNumber,247 imageUrls: rec.imageUrls,248 rawTitle: rec.rawTitle,249 description: rec.description,250 listedAt: rec.listedAt,251 endsAt: rec.endsAt,252 availability: rec.availability,253 bidCount: rec.bidCount,254 firstSeenAt: now,255 lastSeenAt: now,256 confidence,257 dataQuality: quality,258 flags: confidence < 0.6 ? ['low_identification_confidence'] : [],259 });260 await db().insert(listingEvents).values({ id: newId('event'), listingId: id, eventType: 'new', newPrice: rec.price, currency: rec.currency, occurredAt: now });261 if (rec.imageUrls.length) await storeImages(r.assetId, rec.imageUrls.slice(0, 3), rec.sourceId, { listingId: id });262 await trackCert(rec, 'listing', id, r.assetId, variant.id, priceUsd, rec.price, rec.currency, now);263 return { assetId: r.assetId, variantId: variant.id, targetId: id, method: r.method, confidence, status: 'applied', event: 'listing_created' };264 }265 const priceChanged = rec.price !== null && existing.price !== null && Math.abs(Number(existing.price) - rec.price) > 0.009;266 const relisted = existing.availability !== 'available' && rec.availability === 'available';267 const ended = existing.availability === 'available' && rec.availability !== 'available';268 await db()269 .update(listings)270 .set({271 price: rec.price,272 currency: rec.currency,273 priceUsd,274 availability: rec.availability,275 bidCount: rec.bidCount,276 endsAt: rec.endsAt ?? sql`${listings.endsAt}`,277 lastSeenAt: now,278 priceChangedAt: priceChanged ? now : sql`${listings.priceChangedAt}`,279 quantity: rec.quantity,280 dataQuality: quality,281 updatedAt: new Date(),282 })283 .where(eq(listings.id, existing.id));284 const ev = priceChanged ? 'price_changed' : relisted ? 'relisted' : ended ? (rec.availability === 'sold' ? 'sold' : rec.availability === 'ended' ? 'auction_ended' : 'removed') : null;285 if (ev) await db().insert(listingEvents).values({ id: newId('event'), listingId: existing.id, eventType: ev, oldPrice: existing.price, newPrice: rec.price, currency: rec.currency, occurredAt: now });286 return { assetId: r.assetId, variantId: variant.id, targetId: existing.id, method: r.method, confidence, status: 'applied', event: ev ? 'listing_updated' : undefined };287}288289async function applyObservation(rec: NormalizedPriceObservation): Promise<ApplyResult> {290 const grade = refineGrade(rec.grade, rec.rawTitle);291 const r = await resolveAsset(rec.attributes, { imageUrls: rec.imageUrls });292 if (!r) return { assetId: null, variantId: null, targetId: null, method: null, confidence: null, status: 'unmatched' };293 const condition = normCond(rec.attributes.categorySlug, rec.condition.conditionRaw, rec.condition.condition);294 const variant = await ensureVariant(r.assetId, grade, { condition, completeness: rec.condition.completeness }, rec.attributes.size);295 const fx = await toUsd(rec.price, rec.currency, rec.observationDate);296 if (!fx) throw new FxMissingError(rec.currency, rec.observationDate);297 const dedupeKey = sha256(['obs', rec.sourceId, rec.externalId ?? rec.sourceUrl, rec.priceKind, variant.key, toDateOnly(rec.observationDate)].join('|'));298 const id = newId('valuation');299 const inserted = await db()300 .insert(priceObservations)301 .values({ id, assetId: r.assetId, variantId: variant.id, sourceId: rec.sourceId, connectorId: rec.connectorId, sourceUrl: rec.sourceUrl, priceKind: rec.priceKind, price: rec.price, currency: rec.currency, priceUsd: fx.usd, observationDate: toDateOnly(rec.observationDate), sampleSize: rec.sampleSize, dedupeKey })302 .onConflictDoNothing({ target: priceObservations.dedupeKey })303 .returning({ id: priceObservations.id });304 if (inserted.length === 0) return { assetId: r.assetId, variantId: variant.id, targetId: null, method: r.method, confidence: r.confidence, status: 'duplicate' };305 return { assetId: r.assetId, variantId: variant.id, targetId: id, method: r.method, confidence: r.confidence, status: 'applied' };306}307308async function applyAuctionLot(rec: NormalizedAuctionLot): Promise<ApplyResult> {309 const grade = refineGrade(rec.grade, rec.rawTitle);310 const r = await resolveAsset(enrichAttributesFromTitle(rec.attributes, rec.rawTitle), { imageUrls: rec.imageUrls, allowCreate: rec.confidence >= 0.7 });311 const auctionUrl = (rec.attributes.metadata?.auctionUrl as string | undefined) ?? `${new URL(rec.sourceUrl).origin}#${rec.auctionHouse}:${rec.auctionName ?? 'auction'}`;312 const auctionId = sha256(auctionUrl).slice(0, 24);313 await db()314 .insert(auctions)315 .values({ id: `auc_${auctionId}`, sourceId: rec.sourceId, auctionHouse: rec.auctionHouse, name: rec.auctionName ?? rec.auctionHouse, url: auctionUrl, startsAt: rec.startsAt, endsAt: rec.endsAt, location: rec.location, categorySlugs: [rec.attributes.categorySlug], status: rec.status === 'unknown' ? 'upcoming' : rec.status, currency: rec.currency })316 .onConflictDoUpdate({ target: auctions.url, set: { endsAt: sql`coalesce(excluded.ends_at, ${auctions.endsAt})`, status: sql`excluded.status`, updatedAt: new Date() } });317 let variantId: string | null = null;318 if (r) variantId = (await ensureVariant(r.assetId, grade, { condition: normCond(rec.attributes.categorySlug, rec.condition.conditionRaw, rec.condition.condition), completeness: rec.condition.completeness }, rec.attributes.size)).id;319 const id = `lot_${sha256(rec.sourceUrl).slice(0, 24)}`;320 await db()321 .insert(auctionLots)322 .values({ id, auctionId: `auc_${auctionId}`, assetId: r?.assetId ?? null, variantId, sourceId: rec.sourceId, lotNumber: rec.lotNumber, title: rec.rawTitle, url: rec.sourceUrl, estimateLow: rec.estimateLow, estimateHigh: rec.estimateHigh, currentBid: rec.currentBid, currency: rec.currency, startsAt: rec.startsAt, endsAt: rec.endsAt, status: rec.status === 'unknown' ? 'upcoming' : rec.status, imageUrls: rec.imageUrls, grader: grade.grader, grade: grade.grade })323 .onConflictDoUpdate({ target: auctionLots.url, set: { currentBid: sql`excluded.current_bid`, status: sql`excluded.status`, endsAt: sql`coalesce(excluded.ends_at, ${auctionLots.endsAt})`, assetId: sql`coalesce(${auctionLots.assetId}, excluded.asset_id)`, updatedAt: new Date() } });324 await trackCert(rec, 'auction_lot', id, r?.assetId ?? null, variantId, null, rec.currentBid, rec.currency, rec.observedAt);325 return { assetId: r?.assetId ?? null, variantId, targetId: id, method: r?.method ?? null, confidence: r?.confidence ?? null, status: r ? 'applied' : 'unmatched' };326}327328async function applyPopulation(rec: NormalizedPopulationReport): Promise<ApplyResult> {329 const r = await resolveAsset(rec.attributes, {});330 if (!r) return { assetId: null, variantId: null, targetId: null, method: null, confidence: null, status: 'unmatched' };331 const id = newId('event');332 const inserted = await db()333 .insert(populationReports)334 .values({ id, assetId: r.assetId, grader: rec.grader, sourceId: rec.sourceId, sourceUrl: rec.sourceUrl, reportDate: toDateOnly(rec.reportDate), total: rec.total, byGrade: rec.byGrade })335 .onConflictDoNothing()336 .returning({ id: populationReports.id });337 return { assetId: r.assetId, variantId: null, targetId: inserted[0]?.id ?? null, method: r.method, confidence: r.confidence, status: inserted.length ? 'applied' : 'duplicate' };338}339340async function applyNews(rec: NormalizedNewsItem): Promise<ApplyResult> {341 const id = newId('news');342 const inserted = await db()343 .insert(news)344 .values({ id, sourceId: rec.sourceId, url: rec.sourceUrl, title: rec.title, summary: rec.summary, publishedAt: rec.publishedAt, categorySlugs: rec.categorySlugs, imageUrl: rec.imageUrl, fetchedAt: new Date() })345 .onConflictDoNothing({ target: news.url })346 .returning({ id: news.id });347 return { assetId: null, variantId: null, targetId: inserted[0]?.id ?? null, method: null, confidence: null, status: inserted.length ? 'applied' : 'duplicate' };348}349