// Auteur : Simon-Pierre Boucher — contact@spboucher.ai /** * Persistance de l'analyse IA des annonces (SERVEUR) : analyses versionnées * (jamais écrasées), images, overrides, conflits, vecteurs, usage API. * Colonnes additionnelles ajoutées par `ensureAiSchema` (ALTER tolérant) — * db.ts n'est pas modifié. */ import type Database from "better-sqlite3"; import { getCostDb } from "../db"; import type { Conflict } from "./merge"; let ensured = false; export function ensureAiSchema(d: Database.Database = getCostDb()): Database.Database { if (ensured) return d; const cols: [string, string][] = [ ["listing_ai_analyses", "merged_json TEXT"], ["listing_ai_analyses", "mapping_json TEXT"], ["listing_ai_analyses", "usage_json TEXT"], ["listing_ai_analyses", "issues_json TEXT"], ["listing_ai_analyses", "images_json TEXT"], ["listing_ai_analyses", "listing_json TEXT"], ["listing_ai_analyses", "combined_json TEXT"], ["listing_ai_analyses", "started_at TEXT"], ["listing_ai_analyses", "merged_base_json TEXT"], ["listing_images", "analysis_id TEXT"], ]; for (const [t, c] of cols) { try { d.exec(`ALTER TABLE ${t} ADD COLUMN ${c}`); } catch { /* déjà présent */ } } d.exec("CREATE INDEX IF NOT EXISTS idx_ai_usage_created ON ai_usage(created_at)"); ensured = true; return d; } export type AnalysisStatus = "queued" | "fetching_images" | "analyzing" | "validating" | "embedding" | "mapping_assemblies" | "pricing" | "completed" | "failed"; export const STAGES: AnalysisStatus[] = ["queued", "fetching_images", "analyzing", "validating", "embedding", "mapping_assemblies", "pricing", "completed"]; export interface AnalysisRow { id: string; listing_uid: string; version: number; model: string; prompt_version: string; schema_version: string; input_hash: string; status: AnalysisStatus; stage: string | null; error: string | null; image_count: number; output_json: string | null; compact_json: string | null; vector_json: string | null; confidence: number | null; estimate_id: string | null; cost_snapshot_date: string | null; created_at: string; completed_at: string | null; merged_json: string | null; merged_base_json: string | null; mapping_json: string | null; usage_json: string | null; issues_json: string | null; images_json: string | null; listing_json: string | null; combined_json: string | null; started_at: string | null; } export function db(): Database.Database { return ensureAiSchema(getCostDb()); } export function findCompleted(listingUid: string, inputHash: string, promptVersion: string, model: string): AnalysisRow | null { return (db().prepare("SELECT * FROM listing_ai_analyses WHERE listing_uid=? AND input_hash=? AND prompt_version=? AND model=? AND status='completed' ORDER BY version DESC LIMIT 1").get(listingUid, inputHash, promptVersion, model) as AnalysisRow | undefined) ?? null; } export function findActive(listingUid: string): AnalysisRow | null { return (db().prepare("SELECT * FROM listing_ai_analyses WHERE listing_uid=? AND status NOT IN ('completed','failed') ORDER BY version DESC LIMIT 1").get(listingUid) as AnalysisRow | undefined) ?? null; } export function latestForListing(listingUid: string): AnalysisRow | null { return (db().prepare("SELECT * FROM listing_ai_analyses WHERE listing_uid=? ORDER BY version DESC LIMIT 1").get(listingUid) as AnalysisRow | undefined) ?? null; } export function latestCompletedForListing(listingUid: string): AnalysisRow | null { return (db().prepare("SELECT * FROM listing_ai_analyses WHERE listing_uid=? AND status='completed' ORDER BY version DESC LIMIT 1").get(listingUid) as AnalysisRow | undefined) ?? null; } export function getAnalysis(id: string): AnalysisRow | null { return (db().prepare("SELECT * FROM listing_ai_analyses WHERE id=?").get(id) as AnalysisRow | undefined) ?? null; } export function listVersions(listingUid: string): { id: string; version: number; status: string; created_at: string; confidence: number | null; model: string; prompt_version: string }[] { return db().prepare("SELECT id, version, status, created_at, confidence, model, prompt_version FROM listing_ai_analyses WHERE listing_uid=? ORDER BY version DESC").all(listingUid) as { id: string; version: number; status: string; created_at: string; confidence: number | null; model: string; prompt_version: string }[]; } /** Analyses complètes par annonce (badges de cartes) — une requête pour une liste d'uids. */ export function completedSummaries(uids: string[]): Map { const out = new Map(); if (!uids.length) return out; const d = db(); const rows = d.prepare(`SELECT a.listing_uid, a.id, a.confidence, e.replacement_cost_new rcn FROM listing_ai_analyses a LEFT JOIN cost_estimates e ON e.id=a.estimate_id WHERE a.status='completed' AND a.listing_uid IN (${uids.map(() => "?").join(",")}) ORDER BY a.version ASC`).all(...uids) as { listing_uid: string; id: string; confidence: number | null; rcn: number | null }[]; for (const r of rows) out.set(r.listing_uid, { analysisId: r.id, rcn: r.rcn, confidence: r.confidence }); return out; } export function createAnalysis(a: { id: string; listingUid: string; model: string; promptVersion: string; schemaVersion: string; inputHash: string; listingJson: string }): AnalysisRow { const d = db(); const v = (d.prepare("SELECT COALESCE(MAX(version),0)+1 v FROM listing_ai_analyses WHERE listing_uid=?").get(a.listingUid) as { v: number }).v; d.prepare("INSERT INTO listing_ai_analyses(id,listing_uid,version,model,prompt_version,schema_version,input_hash,status,stage,listing_json,started_at) VALUES(?,?,?,?,?,?,?,'queued','queued',?,datetime('now'))").run(a.id, a.listingUid, v, a.model, a.promptVersion, a.schemaVersion, a.inputHash, a.listingJson); return getAnalysis(a.id)!; } export function setStage(id: string, status: AnalysisStatus): void { db().prepare("UPDATE listing_ai_analyses SET status=?, stage=? WHERE id=?").run(status, status, id); } export function setFailed(id: string, error: string): void { db().prepare("UPDATE listing_ai_analyses SET status='failed', stage='failed', error=?, completed_at=datetime('now') WHERE id=?").run(error.slice(0, 2000), id); } /** * Analyses orphelines : une analyse tourne dans le processus du serveur ; un * redémarrage (pm2 restart, déploiement) l'interrompt sans transition vers * `failed`, et l'interface resterait sur « analyse en cours » indéfiniment. * Toute analyse non terminale plus vieille que `maxAgeMin` est marquée échouée * (les plus longues observées : ≈ 2 min pour 27-36 photos). Renvoie le nombre corrigé. */ export function failOrphans(maxAgeMin = 20): number { const r = db().prepare( "UPDATE listing_ai_analyses SET status='failed', stage='failed', error=?, completed_at=datetime('now') WHERE status NOT IN ('completed','failed') AND COALESCE(started_at, created_at) < datetime('now', ?)", ).run("Analyse interrompue (redémarrage du serveur pendant le traitement) — relancez l'analyse.", `-${Math.max(1, Math.round(maxAgeMin))} minutes`); return r.changes; } export function saveOutput(id: string, o: { outputJson: string; imageCount: number; usageJson: string; issuesJson: string; imagesJson: string }): void { db().prepare("UPDATE listing_ai_analyses SET output_json=?, image_count=?, usage_json=?, issues_json=?, images_json=? WHERE id=?").run(o.outputJson, o.imageCount, o.usageJson, o.issuesJson, o.imagesJson, id); } /** Faits fusionnés de BASE (avant corrections manuelles) + vecteur. */ export function saveMergedAndVector(id: string, mergedJson: string, vectorJson: string): void { db().prepare("UPDATE listing_ai_analyses SET merged_base_json=?, merged_json=?, vector_json=? WHERE id=?").run(mergedJson, mergedJson, vectorJson, id); } export function saveMergedBase(id: string, mergedJson: string): void { db().prepare("UPDATE listing_ai_analyses SET merged_base_json=? WHERE id=?").run(mergedJson, id); } /** Faits EFFECTIFS (corrections appliquées) tels que présentés et chiffrés. */ export function saveMergedEffective(id: string, mergedJson: string): void { db().prepare("UPDATE listing_ai_analyses SET merged_json=? WHERE id=?").run(mergedJson, id); } export function saveMappingAndEstimate(id: string, o: { mappingJson: string; compactJson: string; estimateId: string; snapshotDate: string; combinedJson: string; confidence: number; completed: boolean }): void { db().prepare(`UPDATE listing_ai_analyses SET mapping_json=?, compact_json=?, estimate_id=?, cost_snapshot_date=?, combined_json=?, confidence=?${o.completed ? ", status='completed', stage='completed', completed_at=datetime('now')" : ""} WHERE id=?`) .run(o.mappingJson, o.compactJson, o.estimateId, o.snapshotDate, o.combinedJson, o.confidence, id); } /** Hashes déjà connus des photos d'une annonce (URL → hash) — permet de recalculer l'empreinte sans retélécharger. */ export function knownImageHashes(listingUid: string): Map { return new Map((db().prepare("SELECT source_url, hash FROM listing_images WHERE listing_uid=? AND hash IS NOT NULL").all(listingUid) as { source_url: string; hash: string }[]).map((r) => [r.source_url, r.hash])); } export function saveImages(listingUid: string, analysisId: string, imgs: { id: string; sourceUrl: string; position: number; width: number | null; height: number | null; bytes: number; hash: string; mediaType: string; roomHint: string; qualityScore: number; selected: boolean }[]): void { const d = db(); const up = d.prepare(`INSERT INTO listing_images(listing_uid,photo_id,source_url,position,width,height,bytes,hash,media_type,room_guess,quality_score,selected_for_ai,fetched_at,analysis_id) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,datetime('now'),?) ON CONFLICT(listing_uid, source_url) DO UPDATE SET photo_id=excluded.photo_id, width=excluded.width, height=excluded.height, bytes=excluded.bytes, hash=excluded.hash, media_type=excluded.media_type, room_guess=excluded.room_guess, quality_score=excluded.quality_score, selected_for_ai=excluded.selected_for_ai, fetched_at=excluded.fetched_at, analysis_id=excluded.analysis_id`); const tx = d.transaction(() => { for (const i of imgs) up.run(listingUid, i.id, i.sourceUrl, i.position, i.width, i.height, i.bytes, i.hash, i.mediaType, i.roomHint, i.qualityScore, i.selected ? 1 : 0, analysisId); }); tx(); } export function saveConflicts(analysisId: string, conflicts: Conflict[]): void { const d = db(); d.prepare("DELETE FROM analysis_conflicts WHERE analysis_id=?").run(analysisId); const ins = d.prepare("INSERT INTO analysis_conflicts(analysis_id,field,source_a,value_a,source_b,value_b,severity,status) VALUES(?,?,?,?,?,?,?,'open')"); const tx = d.transaction(() => { for (const c of conflicts) ins.run(analysisId, c.field, c.sourceA, c.valueA, c.sourceB, c.valueB, c.severity); }); tx(); } export function listConflicts(analysisId: string): (Conflict & { id: number; status: string })[] { return (db().prepare("SELECT id, field, source_a, value_a, source_b, value_b, severity, status FROM analysis_conflicts WHERE analysis_id=?").all(analysisId) as { id: number; field: string; source_a: string; value_a: string; source_b: string; value_b: string; severity: string; status: string }[]) .map((r) => ({ id: r.id, field: r.field, sourceA: r.source_a as Conflict["sourceA"], valueA: r.value_a, sourceB: r.source_b as Conflict["sourceB"], valueB: r.value_b, severity: r.severity as Conflict["severity"], status: r.status })); } export function saveEmbedding(listingUid: string, analysisId: string, model: string, vector: number[], text: string): void { db().prepare("INSERT INTO property_embeddings(listing_uid,analysis_id,embedding_model,embedding_dimension,embedding_json,canonical_text) VALUES(?,?,?,?,?,?) ON CONFLICT(analysis_id, embedding_model) DO UPDATE SET embedding_json=excluded.embedding_json, canonical_text=excluded.canonical_text, embedding_dimension=excluded.embedding_dimension") .run(listingUid, analysisId, model, vector.length, JSON.stringify(vector), text); } export function allEmbeddings(model: string): { listingUid: string; analysisId: string; vector: number[]; text: string | null }[] { return (db().prepare(`SELECT e.listing_uid, e.analysis_id, e.embedding_json, e.canonical_text FROM property_embeddings e JOIN listing_ai_analyses a ON a.id=e.analysis_id WHERE e.embedding_model=? AND a.status='completed' AND a.version = (SELECT MAX(version) FROM listing_ai_analyses b WHERE b.listing_uid=a.listing_uid AND b.status='completed')`).all(model) as { listing_uid: string; analysis_id: string; embedding_json: string; canonical_text: string | null }[]) .map((r) => ({ listingUid: r.listing_uid, analysisId: r.analysis_id, vector: JSON.parse(r.embedding_json) as number[], text: r.canonical_text })); } export function addOverride(analysisId: string, fieldPath: string, original: unknown, value: unknown, userId: string | null): void { db().prepare("INSERT INTO listing_ai_overrides(analysis_id,field_path,original_value,new_value,user_id) VALUES(?,?,?,?,?)").run(analysisId, fieldPath, JSON.stringify(original ?? null), JSON.stringify(value ?? null), userId); } export function listOverrides(analysisId: string): { fieldPath: string; original: unknown; value: unknown; createdAt: string }[] { const rows = db().prepare("SELECT field_path, original_value, new_value, created_at FROM listing_ai_overrides WHERE analysis_id=? ORDER BY id").all(analysisId) as { field_path: string; original_value: string; new_value: string; created_at: string }[]; // la dernière correction d'un champ gagne const map = new Map(); for (const r of rows) map.set(r.field_path, { fieldPath: r.field_path, original: map.get(r.field_path)?.original ?? JSON.parse(r.original_value), value: JSON.parse(r.new_value), createdAt: r.created_at }); return [...map.values()]; } export function recordUsage(u: { analysisId: string | null; purpose: string; model: string; images: number; inputTokens: number; outputTokens: number; cacheRead: number; costUsd: number; latencyMs: number }): void { db().prepare("INSERT INTO ai_usage(analysis_id,purpose,model,input_images,input_tokens,output_tokens,cache_read_tokens,estimated_cost_usd,latency_ms) VALUES(?,?,?,?,?,?,?,?,?)").run(u.analysisId, u.purpose, u.model, u.images, u.inputTokens, u.outputTokens, u.cacheRead, u.costUsd, u.latencyMs); } export function usageToday(): { count: number; costUsd: number } { const r = db().prepare("SELECT COUNT(*) n, COALESCE(SUM(estimated_cost_usd),0) c FROM ai_usage WHERE created_at >= date('now')").get() as { n: number; c: number }; return { count: r.n, costUsd: r.c }; } export function analysesStartedToday(): number { return (db().prepare("SELECT COUNT(*) n FROM listing_ai_analyses WHERE created_at >= date('now')").get() as { n: number }).n; }