import crypto from "node:crypto"; import sharp from "sharp"; import { z } from "zod"; import type { Db } from "../db.js"; import { tx, nowIso } from "../db.js"; import type { AppConfig } from "../config.js"; import { newId } from "./ids.js"; import { audit } from "./audit.js"; import { enqueueJob, getJob, finishJob } from "./jobs.js"; import { upsertFabric, upsertShade } from "./fabric.js"; import { ingestUpload } from "../media/ingest.js"; import { ensureImportSnapshot, markImportStored, handleImportFailure, getSourceImport, SourceImportError, type SourceImportRow, } from "./sourceImports.js"; import { normalizeSourceUrl, fetchSource, type SourceFetcher } from "../media/sourceFetch.js"; import { parseFabricCollectionPage, detectFabricVendor, FabricParseError } from "../import/fabricParser.js"; export class FabricImportError extends Error { constructor(public readonly statusCode: number, message: string) { super(message); } } export const enqueueFabricImportSchema = z.object({ url: z.string().min(1).max(2048), fabricCode: z.string().regex(/^[\w-]{1,40}$/).optional(), idempotencyKey: z.string().regex(/^[\w-]{1,128}$/), }).strict(); const hash = (value: string) => crypto.createHash("sha256").update(value).digest("hex"); /** Średni kolor próbki → hex "#rrggbb" (wzorce: scripts/import-fuji.mts rgbHex). */ async function dominantHex(buffer: Buffer): Promise { try { const st = await sharp(buffer).stats(); const [r, g, b] = st.channels.map((c) => Math.round(c.mean)); return "#" + [r, g, b].map((v) => v!.toString(16).padStart(2, "0")).join(""); } catch { return null; } } /** * Zleca import tkaniny z URL strony kolekcji dostawcy (SIC lub Davis). * Atomowo: wiersz source_imports (kind fabric_page) + job import_fabric. * Katalog po imporcie jest zawsze unverified — potwierdza człowiek. */ export function enqueueFabricImport( db: Db, input: { url: string; fabricCode?: string | undefined; idempotencyKey: string; actor: string }, ) { const normalized = normalizeSourceUrl(input.url, "fabric_page"); const fingerprint = hash(JSON.stringify([normalized.requestedUrl, "fabric_page", input.fabricCode ?? null])); const requestKey = hash(JSON.stringify([input.actor, "fabric", input.idempotencyKey])); return tx(db, () => { const prior = db.prepare("SELECT id, request_fingerprint FROM source_imports WHERE request_key = ?").get(requestKey) as { id: string; request_fingerprint: string } | undefined; if (prior) { if (prior.request_fingerprint !== fingerprint) throw new FabricImportError(409, "Ten klucz idempotencji został już użyty z innymi danymi."); return { import: getSourceImport(db, prior.id)!, created: false }; } const pending = db.prepare("SELECT COUNT(*) n FROM source_imports i JOIN jobs j ON j.id = i.job_id WHERE j.status IN ('queued','running','waiting')").get() as { n: number }; if (pending.n >= 20) throw new SourceImportError(429, "Kolejka importu jest pełna. Poczekaj na zakończenie istniejących zadań."); const job = enqueueJob(db, { type: "import_fabric", payload: { fabricCode: input.fabricCode ?? null }, provider: "source-fetch", idempotencyKey: `fabric-import:${requestKey}`, maxAttempts: 3, }); const importId = newId("imp"); const now = nowIso(); db.prepare(`INSERT INTO source_imports (id,job_id,product_id,kind,requested_url,fetch_url,fragment,request_key,request_fingerprint,created_by,created_at,updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?)`) .run(importId, job.id, null, "fabric_page", normalized.requestedUrl, normalized.fetchUrl, normalized.fragment, requestKey, fingerprint, input.actor, now, now); audit(db, { actor: input.actor, action: "fabric_import.queued", entity: "source_imports", entityId: importId, detail: { jobId: job.id, url: normalized.requestedUrl, fabricCode: input.fabricCode ?? null } }); return { import: getSourceImport(db, importId)!, created: true }; }); } /** * Job import_fabric: snapshot strony kolekcji → parser wg dostawcy (SIC/Davis) * → fabrics + fabric_shades (unverified) + próbki jako source_assets * (fabric_reference) + karta PDF (source_document) podpięta jako card_path. * Reimport nie obniża statusów confirmed. Błąd parsera = error bez retry. */ export async function runFabricImport( db: Db, cfg: AppConfig, jobId: string, fetcher: SourceFetcher = fetchSource, stopSignal?: AbortSignal, ): Promise { const row = db.prepare( `SELECT i.*, j.status AS job_status, j.error AS job_error, j.attempts FROM source_imports i JOIN jobs j ON j.id = i.job_id WHERE i.job_id = ?` ).get(jobId) as unknown as SourceImportRow | undefined; if (!row || row.job_status !== "running") return; const job = getJob(db, jobId); const payload = JSON.parse(job?.payload_json ?? "{}") as { fabricCode?: string | null }; const cancel = new AbortController(); const signal = stopSignal ? AbortSignal.any([stopSignal, cancel.signal]) : cancel.signal; const timer = setInterval(() => { if (getJob(db, jobId)?.status === "cancelled") cancel.abort(); }, 100); const warnings: string[] = []; try { const snap = await ensureImportSnapshot(db, cfg, row, jobId, fetcher, signal); if (!snap) return; signal.throwIfAborted(); if (!["text/html", "application/xhtml+xml"].includes(snap.manifest.mime)) { throw new FabricImportError(422, "Snapshot kolekcji nie jest dokumentem HTML."); } const parsed = parseFabricCollectionPage(snap.buffer.toString("utf8"), row.requested_url); const vendor = detectFabricVendor(row.requested_url); signal.throwIfAborted(); if (getJob(db, jobId)?.status === "cancelled") return; markImportStored(db, row, snap.manifest, snap.recovered); const code = payload.fabricCode ?? parsed.code; const existingFabric = db.prepare("SELECT id, params_json, verification FROM fabrics WHERE code = ?").get(code) as { id: string; params_json: string; verification: string } | undefined; const fabricConfirmed = existingFabric?.verification === "confirmed"; const fabric = upsertFabric(db, { code, name: parsed.name.toUpperCase(), producer: parsed.producer, // params tylko przy tworzeniu — reimport nie nadpisuje zweryfikowanych parametrów ...(existingFabric ? {} : { params: { ...parsed.params, source: new URL(row.requested_url).hostname.replace(/^www\./, ""), collectionUrl: row.requested_url } }), }); db.prepare("UPDATE source_imports SET fabric_id = ?, updated_at = ? WHERE id = ?").run(fabric.id, nowIso(), row.id); // Karta PDF — błąd pobrania nie przerywa importu odcieni (warning). let cardStored = false; if (parsed.cardUrl) { try { const card = await fetcher(parsed.cardUrl, "fabric_card", signal); const doc = await ingestUpload(db, cfg, { buffer: card.buffer, originalName: `${code}-${vendor}-product-profile.pdf`, productId: null, role: "source_document", note: `Karta produktu ${parsed.name.toUpperCase()} (${parsed.producer}) — źródło parametrów tkaniny. Import ${row.id}.`, uploadedBy: `import:${row.id}`, sourceUrl: card.finalUrl, fetchedAt: card.fetchedAt, }); // potwierdzonej tkaninie nie podmieniamy karty — kuratorskie dane zostają upsertFabric(db, { code, name: parsed.name.toUpperCase(), producer: parsed.producer, cardPath: fabricConfirmed ? null : doc.path }); cardStored = true; } catch (err) { warnings.push(`karta PDF: ${err instanceof Error ? err.message : String(err)}`); } } else { warnings.push("brak linku do karty PDF na stronie kolekcji"); } // Próbki odcieni — każda przez safe fetcher (host allowlist, limity). let swatchesStored = 0; let skippedConfirmed = 0; for (const shade of parsed.shades) { signal.throwIfAborted(); if (getJob(db, jobId)?.status === "cancelled") return; const existing = db.prepare("SELECT verification FROM fabric_shades WHERE fabric_id = ? AND code = ?") .get(fabric.id, shade.code) as { verification: string } | undefined; const verification = existing?.verification === "confirmed" ? "confirmed" : "unverified"; if (existing?.verification === "confirmed") skippedConfirmed++; let officialImagePath: string | null = null; let colorHex: string | null = null; try { const img = await fetcher(shade.imageUrl, "fabric_sample", signal); colorHex = await dominantHex(img.buffer); const out = await ingestUpload(db, cfg, { buffer: img.buffer, originalName: `${code}-${shade.code}${shade.name ? `-${shade.name.replace(/\s+/g, "-")}` : ""}.jpg`, productId: null, role: "fabric_reference", note: `Oficjalna próbka producenta — ${parsed.name.toUpperCase()} ${shade.code}${shade.name ? ` (${shade.name})` : ""}. Import ${vendor} ${row.id}.`, uploadedBy: `import:${row.id}`, sourceUrl: img.finalUrl, fetchedAt: img.fetchedAt, }); officialImagePath = out.id; if (!out.duplicate) swatchesStored++; } catch (err) { warnings.push(`odcień ${shade.code}: ${err instanceof Error ? err.message : String(err)}`); } upsertShade(db, fabric.id, { code: shade.code, name: shade.name, colorHex, // potwierdzonemu odcieniowi nie podmieniamy próbki officialImagePath: existing?.verification === "confirmed" ? null : officialImagePath, verification, }); } tx(db, () => { finishJob(db, jobId, { ok: true, result: { importId: row.id, fabricId: fabric.id, code, collection: parsed.name, sha256: snap.manifest.sha256, shadesParsed: parsed.shades.length, swatchesStored, cardStored, skippedConfirmed, warnings, }, }); audit(db, { actor: row.created_by, action: "fabric_import.done", entity: "fabrics", entityId: fabric.id, detail: { importId: row.id, code, shades: parsed.shades.length, swatchesStored, cardStored, warnings } }); }); } catch (err) { if (err instanceof FabricParseError || err instanceof FabricImportError) { if (getJob(db, jobId)?.status === "cancelled" || stopSignal?.aborted) { handleImportFailure(db, row, jobId, err, stopSignal); return; } tx(db, () => { finishJob(db, jobId, { ok: false, error: err.message, retryable: false }); db.prepare("UPDATE source_imports SET phase='failed', error_code=?, updated_at=? WHERE id=?") .run(err instanceof FabricParseError ? `parse_${err.code}` : "import_error", nowIso(), row.id); audit(db, { actor: row.created_by, action: "fabric_import.failed", entity: "source_imports", entityId: row.id, detail: { error: err.message } }); }); return; } handleImportFailure(db, row, jobId, err, stopSignal); } finally { clearInterval(timer); } }