SPB Git forge

spb/rareindex

Public
54commits 1branches 0releases
7.1 MBsize
maindefault branch
10 days agolast push
TypeScript 61.9% HTML 37.2% SQL 0.7%
21.6 KB · 349 lines typescript
Raw Blame History
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