import { and, eq, sql } from 'drizzle-orm'; import { auctionLots, auctions, images, listingEvents, listings, news, populationReports, priceObservations, sales, sources } from '@rareindex/database'; import { normalizeCondition, parseGradeFromTitle } from '@rareindex/taxonomy'; import { allInPrice, feeScheduleFor } from '@rareindex/valuation'; import type { Grade } from '@rareindex/shared'; import { newId, sha256, toDateOnly, type NormalizedAuctionLot, type NormalizedCatalogItem, type NormalizedListing, type NormalizedNewsItem, type NormalizedPopulationReport, type NormalizedPriceObservation, type NormalizedRecord, type NormalizedSale } from '@rareindex/shared'; import { db } from '../lib/db.ts'; import { toUsd, usdRateFor } from '../lib/fx.ts'; import { audit } from '../lib/audit.ts'; import { enrichAttributesFromTitle, ensureVariant, resolveAsset, type Resolution } from './resolver.ts'; import { recordCertificate } from './certificates.ts'; /** Record a slab certification sighting when a record carries one (SPEC §22); never throws into the writer. */ async 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 { const g = refineGrade(rec.grade, rec.rawTitle); if (!g.certificationNumber || !g.grader) return; try { 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 }); } catch { /* cert tracking is best-effort; the canonical write already succeeded */ } } export class FxMissingError extends Error { constructor(public currency: string, public date: Date) { super(`fx_missing:${currency}:${toDateOnly(date)}`); } } export interface ApplyResult { assetId: string | null; variantId: string | null; targetId: string | null; method: Resolution['method'] | null; confidence: number | null; status: 'applied' | 'unmatched' | 'rejected' | 'duplicate'; reason?: string; event?: 'sale_detected' | 'listing_created' | 'listing_updated' | 'entity_created'; } const sourceMetaCache = new Map(); async function sourceMeta(sourceId: string): Promise<{ name: string; sourceType: string } | null> { const hit = sourceMetaCache.get(sourceId); if (hit) return hit; const [row] = await db().select({ name: sources.name, sourceType: sources.sourceType }).from(sources).where(eq(sources.id, sourceId)).limit(1); if (!row) return null; sourceMetaCache.set(sourceId, row); return row; } /** * Buyer-pays price of a sale (§35): price_usd + the house's buyer premium when the record is * hammer-only. `fxToScheduleCurrency` converts the native amount into the schedule's currency so * marginal tiers apply at the right thresholds. Returns nulls when the basis cannot be established. */ export 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 }> { const meta = await sourceMeta(rec.sourceId); const house = rec.auctionHouse ?? meta?.name ?? rec.connectorId; const schedule = feeScheduleFor(house) ?? feeScheduleFor(rec.connectorId); let fxToScheduleCurrency = 1; if (schedule && schedule.currency !== 'USD') { // price is already USD here; tiers are in the schedule currency → USD × (schedule units per USD) const r = await usdRateFor(schedule.currency, rec.saleDate); if (r && r.rate > 0) fxToScheduleCurrency = r.rate; } const a = allInPrice({ price: rec.priceUsd, buyerPremiumIncluded: rec.buyerPremiumIncluded, house: schedule ? schedule.id : house, saleType: rec.saleType, sourceType: meta?.sourceType ?? null, fxToScheduleCurrency }); return { allInUsd: a.allIn, feeBasis: a.basis, buyerPremiumRate: a.rate }; } const trustCache = new Map(); async function sourceTrust(sourceId: string): Promise { const hit = trustCache.get(sourceId); if (hit !== undefined) return hit; const [row] = await db().select({ t: sources.trustScore }).from(sources).where(eq(sources.id, sourceId)).limit(1); const t = row ? Number(row.t) : 0.5; trustCache.set(sourceId, t); return t; } /** Data quality 0–100 (§150): completeness × source trust × identification confidence. */ function dataQuality(fields: Record, trust: number, confidence: number, extra = 0): number { const keys = Object.keys(fields); const present = keys.filter((k) => fields[k] !== null && fields[k] !== undefined && fields[k] !== '' && !(Array.isArray(fields[k]) && (fields[k] as unknown[]).length === 0)).length; const completeness = keys.length ? present / keys.length : 0; return Math.round(Math.max(0, Math.min(100, (0.4 * completeness + 0.3 * trust + 0.3 * confidence + extra) * 100))); } function normCond(categorySlug: string, raw: string | null | undefined, already: string | null | undefined): string | null { return already ?? normalizeCondition(categorySlug, raw); } async function storeImages(assetId: string, urls: string[], sourceId: string, ref: { listingId?: string; saleId?: string }): Promise { if (!urls.length) return; 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' })); await db().insert(images).values(rows).onConflictDoNothing({ target: images.url }); } export async function applyRecord(rec: NormalizedRecord): Promise { switch (rec.kind) { case 'catalog_item': return applyCatalog(rec); case 'sale': return applySale(rec); case 'listing': return applyListing(rec); case 'price_observation': return applyObservation(rec); case 'auction_lot': return applyAuctionLot(rec); case 'population_report': return applyPopulation(rec); case 'news_item': return applyNews(rec); } } /** * Grade refinement (§18, §84): when a connector did not read the slab (or read the grader but not the * grade), parse the raw title. A grader with no readable grade stays "grader · grade unknown" — it must * never fall into the raw/ungraded variant, which is what compared PSA slabs against raw valuations. * An explicit `raw` from the connector is trusted over the title. */ export function refineGrade(grade: Grade, rawTitle: string | null | undefined): Grade { if (!rawTitle || grade.grader === 'raw' || (grade.grader && grade.grade)) return grade; const parsed = parseGradeFromTitle(rawTitle); if (!parsed.grader || parsed.grader === 'raw') return grade; if (grade.grader && parsed.grader !== grade.grader) return grade; // never contradict the connector's grader // A grade recorded without a grader is a condition-scale value (e.g. a price-guide tab), not a slab // grade: when the title reveals the slab, the title's grade wins. const g = grade.grader ? (grade.grade ?? parsed.grade) : parsed.grade; return { ...grade, grader: grade.grader ?? parsed.grader, grade: g, qualifier: grade.qualifier ?? parsed.qualifier }; } async function applyCatalog(rec: NormalizedCatalogItem): Promise { const r = await resolveAsset(rec.attributes, { enrich: true, imageUrls: rec.imageUrls, description: rec.description }); if (!r) return { assetId: null, variantId: null, targetId: null, method: null, confidence: null, status: 'unmatched' }; if (rec.imageUrls.length) await storeImages(r.assetId, rec.imageUrls.slice(0, 1), rec.sourceId, {}); return { assetId: r.assetId, variantId: null, targetId: r.assetId, method: r.method, confidence: r.confidence, status: 'applied', event: r.created ? 'entity_created' : undefined }; } async function applySale(rec: NormalizedSale): Promise { const grade = refineGrade(rec.grade, rec.rawTitle); const r = await resolveAsset(enrichAttributesFromTitle(rec.attributes, rec.rawTitle), { imageUrls: rec.imageUrls }); if (!r) return { assetId: null, variantId: null, targetId: null, method: null, confidence: null, status: 'unmatched' }; const condition = normCond(rec.attributes.categorySlug, rec.condition.conditionRaw, rec.condition.condition); const variant = await ensureVariant(r.assetId, grade, { condition, completeness: rec.condition.completeness }, rec.attributes.size); const fx = await toUsd(rec.price, rec.currency, rec.saleDate); if (!fx) throw new FxMissingError(rec.currency, rec.saleDate); const dedupeKey = sha256(['sale', rec.sourceId, rec.externalId ?? rec.sourceUrl, toDateOnly(rec.saleDate), rec.price.toFixed(2)].join('|')); const trust = await sourceTrust(rec.sourceId); 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 }); const confidence = Math.min(rec.confidence, r.confidence); const flags: string[] = []; if (rec.isBundle || rec.quantity > 1) flags.push('bundle'); if (rec.price <= 0) flags.push('zero_price'); if (confidence < 0.6) flags.push('low_identification_confidence'); const status = flags.includes('bundle') || flags.includes('zero_price') ? 'excluded' : 'valid'; const id = newId('sale'); const inserted = await db() .insert(sales) .values({ id, assetId: r.assetId, variantId: variant.id, sourceId: rec.sourceId, connectorId: rec.connectorId, sourceUrl: rec.sourceUrl, externalId: rec.externalId, saleType: rec.saleType, saleDate: rec.saleDate, price: rec.price, currency: rec.currency, priceUsd: fx.usd, fxRate: fx.rate, fxDate: fx.fxDate, buyerPremiumIncluded: rec.buyerPremiumIncluded, allInUsd: fees.allInUsd, feeBasis: fees.feeBasis, buyerPremiumRate: fees.buyerPremiumRate, quantity: rec.quantity, isBundle: rec.isBundle, condition, grader: grade.grader, grade: grade.grade, certificationNumber: grade.certificationNumber, location: rec.location, auctionHouse: rec.auctionHouse, lotNumber: rec.lotNumber, imageUrls: rec.imageUrls, rawTitle: rec.rawTitle, confidence, 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), status, flags, dedupeKey, }) .onConflictDoNothing({ target: sales.dedupeKey }) .returning({ id: sales.id }); if (inserted.length === 0) return { assetId: r.assetId, variantId: variant.id, targetId: null, method: r.method, confidence, status: 'duplicate' }; if (status === 'excluded') await audit({ entityType: 'sale', entityId: id, action: 'excluded', reason: flags.join(','), details: { rawTitle: rec.rawTitle, quantity: rec.quantity } }); if (rec.imageUrls.length) await storeImages(r.assetId, rec.imageUrls.slice(0, 3), rec.sourceId, { saleId: id }); await trackCert(rec, 'sale', id, r.assetId, variant.id, fx.usd, rec.price, rec.currency, rec.saleDate); return { assetId: r.assetId, variantId: variant.id, targetId: id, method: r.method, confidence, status: 'applied', event: 'sale_detected' }; } async function applyListing(rec: NormalizedListing): Promise { const grade = refineGrade(rec.grade, rec.rawTitle); const r = await resolveAsset(enrichAttributesFromTitle(rec.attributes, rec.rawTitle), { imageUrls: rec.imageUrls }); if (!r) return { assetId: null, variantId: null, targetId: null, method: null, confidence: null, status: 'unmatched' }; const condition = normCond(rec.attributes.categorySlug, rec.condition.conditionRaw, rec.condition.condition); const variant = await ensureVariant(r.assetId, grade, { condition, completeness: rec.condition.completeness }, rec.attributes.size); let priceUsd: number | null = null; if (rec.price !== null && rec.currency) { const fx = await toUsd(rec.price, rec.currency, rec.observedAt); if (!fx) throw new FxMissingError(rec.currency, rec.observedAt); priceUsd = fx.usd; } const externalId = rec.externalId ?? sha256(rec.sourceUrl); const trust = await sourceTrust(rec.sourceId); const confidence = Math.min(rec.confidence, r.confidence); const now = rec.observedAt; 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); 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); if (!existing) { const id = newId('listing'); await db().insert(listings).values({ id, assetId: r.assetId, variantId: variant.id, sourceId: rec.sourceId, connectorId: rec.connectorId, sourceUrl: rec.sourceUrl, externalId, listingType: rec.listingType, price: rec.price, currency: rec.currency, priceUsd, seller: rec.seller, sellerReputation: rec.sellerReputation, location: rec.location, shippingCost: rec.shippingCost, quantity: rec.quantity, condition, grader: grade.grader, grade: grade.grade, certificationNumber: grade.certificationNumber, imageUrls: rec.imageUrls, rawTitle: rec.rawTitle, description: rec.description, listedAt: rec.listedAt, endsAt: rec.endsAt, availability: rec.availability, bidCount: rec.bidCount, firstSeenAt: now, lastSeenAt: now, confidence, dataQuality: quality, flags: confidence < 0.6 ? ['low_identification_confidence'] : [], }); await db().insert(listingEvents).values({ id: newId('event'), listingId: id, eventType: 'new', newPrice: rec.price, currency: rec.currency, occurredAt: now }); if (rec.imageUrls.length) await storeImages(r.assetId, rec.imageUrls.slice(0, 3), rec.sourceId, { listingId: id }); await trackCert(rec, 'listing', id, r.assetId, variant.id, priceUsd, rec.price, rec.currency, now); return { assetId: r.assetId, variantId: variant.id, targetId: id, method: r.method, confidence, status: 'applied', event: 'listing_created' }; } const priceChanged = rec.price !== null && existing.price !== null && Math.abs(Number(existing.price) - rec.price) > 0.009; const relisted = existing.availability !== 'available' && rec.availability === 'available'; const ended = existing.availability === 'available' && rec.availability !== 'available'; await db() .update(listings) .set({ price: rec.price, currency: rec.currency, priceUsd, availability: rec.availability, bidCount: rec.bidCount, endsAt: rec.endsAt ?? sql`${listings.endsAt}`, lastSeenAt: now, priceChangedAt: priceChanged ? now : sql`${listings.priceChangedAt}`, quantity: rec.quantity, dataQuality: quality, updatedAt: new Date(), }) .where(eq(listings.id, existing.id)); const ev = priceChanged ? 'price_changed' : relisted ? 'relisted' : ended ? (rec.availability === 'sold' ? 'sold' : rec.availability === 'ended' ? 'auction_ended' : 'removed') : null; 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 }); return { assetId: r.assetId, variantId: variant.id, targetId: existing.id, method: r.method, confidence, status: 'applied', event: ev ? 'listing_updated' : undefined }; } async function applyObservation(rec: NormalizedPriceObservation): Promise { const grade = refineGrade(rec.grade, rec.rawTitle); const r = await resolveAsset(rec.attributes, { imageUrls: rec.imageUrls }); if (!r) return { assetId: null, variantId: null, targetId: null, method: null, confidence: null, status: 'unmatched' }; const condition = normCond(rec.attributes.categorySlug, rec.condition.conditionRaw, rec.condition.condition); const variant = await ensureVariant(r.assetId, grade, { condition, completeness: rec.condition.completeness }, rec.attributes.size); const fx = await toUsd(rec.price, rec.currency, rec.observationDate); if (!fx) throw new FxMissingError(rec.currency, rec.observationDate); const dedupeKey = sha256(['obs', rec.sourceId, rec.externalId ?? rec.sourceUrl, rec.priceKind, variant.key, toDateOnly(rec.observationDate)].join('|')); const id = newId('valuation'); const inserted = await db() .insert(priceObservations) .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 }) .onConflictDoNothing({ target: priceObservations.dedupeKey }) .returning({ id: priceObservations.id }); if (inserted.length === 0) return { assetId: r.assetId, variantId: variant.id, targetId: null, method: r.method, confidence: r.confidence, status: 'duplicate' }; return { assetId: r.assetId, variantId: variant.id, targetId: id, method: r.method, confidence: r.confidence, status: 'applied' }; } async function applyAuctionLot(rec: NormalizedAuctionLot): Promise { const grade = refineGrade(rec.grade, rec.rawTitle); const r = await resolveAsset(enrichAttributesFromTitle(rec.attributes, rec.rawTitle), { imageUrls: rec.imageUrls, allowCreate: rec.confidence >= 0.7 }); const auctionUrl = (rec.attributes.metadata?.auctionUrl as string | undefined) ?? `${new URL(rec.sourceUrl).origin}#${rec.auctionHouse}:${rec.auctionName ?? 'auction'}`; const auctionId = sha256(auctionUrl).slice(0, 24); await db() .insert(auctions) .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 }) .onConflictDoUpdate({ target: auctions.url, set: { endsAt: sql`coalesce(excluded.ends_at, ${auctions.endsAt})`, status: sql`excluded.status`, updatedAt: new Date() } }); let variantId: string | null = null; 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; const id = `lot_${sha256(rec.sourceUrl).slice(0, 24)}`; await db() .insert(auctionLots) .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 }) .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() } }); await trackCert(rec, 'auction_lot', id, r?.assetId ?? null, variantId, null, rec.currentBid, rec.currency, rec.observedAt); return { assetId: r?.assetId ?? null, variantId, targetId: id, method: r?.method ?? null, confidence: r?.confidence ?? null, status: r ? 'applied' : 'unmatched' }; } async function applyPopulation(rec: NormalizedPopulationReport): Promise { const r = await resolveAsset(rec.attributes, {}); if (!r) return { assetId: null, variantId: null, targetId: null, method: null, confidence: null, status: 'unmatched' }; const id = newId('event'); const inserted = await db() .insert(populationReports) .values({ id, assetId: r.assetId, grader: rec.grader, sourceId: rec.sourceId, sourceUrl: rec.sourceUrl, reportDate: toDateOnly(rec.reportDate), total: rec.total, byGrade: rec.byGrade }) .onConflictDoNothing() .returning({ id: populationReports.id }); return { assetId: r.assetId, variantId: null, targetId: inserted[0]?.id ?? null, method: r.method, confidence: r.confidence, status: inserted.length ? 'applied' : 'duplicate' }; } async function applyNews(rec: NormalizedNewsItem): Promise { const id = newId('news'); const inserted = await db() .insert(news) .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() }) .onConflictDoNothing({ target: news.url }) .returning({ id: news.id }); return { assetId: null, variantId: null, targetId: inserted[0]?.id ?? null, method: null, confidence: null, status: inserted.length ? 'applied' : 'duplicate' }; }