import fs from "node:fs"; import path from "node:path"; import { z } from "zod"; import type { Db } from "../db.js"; import type { AppConfig } from "../config.js"; import { nowIso, tx } from "../db.js"; import { newId, ID_PREFIX } from "./ids.js"; import { audit } from "./audit.js"; import { contentHash } from "./mediaReview.js"; import { buildOfferPackage } from "./offerPackage.js"; import { getRendition, renditionFile } from "./renditions.js"; import { getAsset } from "./assets.js"; import { enqueueJob, getJob, finishJob, deferJob } from "./jobs.js"; import { previewBaseSchema, translateParentParams, translateVariantParams, materializeMedia, ContractError, type PreviewBase, } from "./baseContract.js"; import { type BaseAdapter, ExportBlocked, ExportRejected, ExportUncertain, } from "../integration/baseAdapter.js"; export class ExportError extends Error { constructor(message: string, readonly statusCode = 400) { super(message); } } const sha256 = z.string().regex(/^[0-9a-f]{64}$/i); export const grantConsentSchema = z.object({ snapshotId: z.string().min(1).max(80), expectedSnapshotSha256: sha256, note: z.string().trim().min(3).max(2000), }).strict(); export const revokeConsentSchema = z.object({ note: z.string().trim().min(3).max(1000), }).strict(); export const dispatchSchema = z.object({ consentId: z.string().min(1).max(80), }).strict(); export const retryItemSchema = z.object({ note: z.string().trim().min(3).max(1000), }).strict(); interface ConsentRow { id: string; product_id: string; snapshot_id: string; snapshot_sha256: string; source_sha256: string; scope_json: string; note: string; granted_by: string; granted_at: string; revoked_at: string | null; revoked_by: string | null; revoke_note: string | null; } interface BatchRow { id: string; product_id: string; consent_id: string; snapshot_id: string; snapshot_sha256: string; job_id: string | null; status: string; created_by: string; created_at: string; updated_at: string; } interface ItemRow { id: string; batch_id: string; product_id: string; kind: "parent" | "variant"; variant_id: string | null; position: number; sku: string | null; method: string; payload_json: string; payload_sha256: string; status: string; external_id: number | null; attempts: number; warnings_json: string | null; last_error: string | null; readback_json: string | null; created_at: string; updated_at: string; } interface SnapshotRow { id: string; product_id: string; source_sha256: string; snapshot_sha256: string; payload_json: string; } function consentRow(db: Db, productId: string, id: string): ConsentRow { const row = db.prepare("SELECT * FROM export_consents WHERE id = ? AND product_id = ?").get(id, productId) as unknown as ConsentRow | undefined; if (!row) throw new ExportError("Nie znaleziono zgody eksportowej dla tego produktu.", 404); return row; } function snapshotRow(db: Db, productId: string, id: string): SnapshotRow { const row = db.prepare("SELECT * FROM offer_snapshots WHERE id = ? AND product_id = ?").get(id, productId) as unknown as SnapshotRow | undefined; if (!row) throw new ExportError("Nie znaleziono snapshotu paczki dla tego produktu.", 404); return row; } function latestDecision(db: Db, snapshotId: string): { decision: string; snapshot_sha256: string } | null { return (db.prepare("SELECT decision, snapshot_sha256 FROM offer_snapshot_decisions WHERE snapshot_id = ? ORDER BY rowid DESC LIMIT 1").get(snapshotId) as { decision: string; snapshot_sha256: string } | undefined) ?? null; } /** * Snapshot dla bieżącego stanu paczki — istniejący dla tego SHA albo * zapisywany w locie. Wysyłka bezpośrednia nie wymaga wcześniejszego * ręcznego zapisywania snapshotu z podglądu. */ function findOrCreateSnapshot(db: Db, productId: string, current: ReturnType, actorId: string): SnapshotRow { const existing = db .prepare("SELECT * FROM offer_snapshots WHERE product_id = ? AND snapshot_sha256 = ?") .get(productId, current.snapshotSha256) as unknown as SnapshotRow | undefined; if (existing) return existing; const version = ((db.prepare("SELECT COALESCE(MAX(version), 0) AS v FROM offer_snapshots WHERE product_id = ?").get(productId) as { v: number }).v) + 1; const id = newId(ID_PREFIX.snapshot); db.prepare( `INSERT INTO offer_snapshots (id, product_id, version, source_sha256, snapshot_sha256, payload_json, created_by, note, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)` ).run(id, productId, version, current.sourceSha256, current.snapshotSha256, JSON.stringify(current), actorId, "Snapshot zapisany automatycznie przy wysyłce do Base.", nowIso()); return db.prepare("SELECT * FROM offer_snapshots WHERE id = ?").get(id) as unknown as SnapshotRow; } /** * Auto-akceptacja snapshotu: brakująca lub nieaktualna decyzja jest * dopisywana jako 'approved' zamiast odrzucać żądanie. Kliknięcie * „Wyślij" przez zalogowanego operatora JEST decyzją człowieka — * wpis trafia do offer_snapshot_decisions i audit_log jak każda inna. */ function ensureSnapshotApproved(db: Db, snapshot: SnapshotRow, actorId: string): void { const decision = latestDecision(db, snapshot.id); if (decision?.decision === "approved" && decision.snapshot_sha256 === snapshot.snapshot_sha256) return; db.prepare( `INSERT INTO offer_snapshot_decisions (id, snapshot_id, snapshot_sha256, decision, decided_by, note, created_at) VALUES (?, ?, ?, 'approved', ?, ?, ?)` ).run(newId(ID_PREFIX.approval), snapshot.id, snapshot.snapshot_sha256, actorId, "Automatyczna akceptacja przy bezpośredniej wysyłce — decyzją jest kliknięcie Wyślij.", nowIso()); audit(db, { actor: actorId, action: "offer.snapshot.approved", entity: "offer_snapshots", entityId: snapshot.id, detail: { productId: snapshot.product_id, sha256: snapshot.snapshot_sha256, auto: true } }); } /** * Rozstrzyga snapshot eksportu bez twardych bramek orkiestratora: * nieaktualne źródła → nowy snapshot w locie, brak akceptacji → auto-decyzja, * blokery w payloadzie → nie odrzucają (raportuje je podgląd, nie wyłącznik). * Zwraca snapshot, do którego należy przywiązać zgodę/wsad. */ function resolveExportSnapshot(db: Db, cfg: AppConfig, productId: string, snapshotId: string, actorId: string): { snapshot: SnapshotRow; base: PreviewBase } { let snapshot = snapshotRow(db, productId, snapshotId); const current = buildOfferPackage(db, cfg, productId); if (current.sourceSha256 !== snapshot.source_sha256) { snapshot = findOrCreateSnapshot(db, productId, current, actorId); } ensureSnapshotApproved(db, snapshot, actorId); const payload = JSON.parse(snapshot.payload_json) as { base?: unknown }; return { snapshot, base: previewBaseSchema.parse(payload.base) }; } function consentScope(base: PreviewBase) { return { target: "base-inventory", method: base.method, inventoryId: base.target.inventoryId, variantIds: base.variants.map((v) => v.variantId), itemCount: base.variants.length + 1, }; } /** Aktywna zgoda dla dokładnego SHA snapshotu albo nowa zgoda auto. */ function findOrCreateConsent(db: Db, productId: string, snapshot: SnapshotRow, base: PreviewBase, actorId: string, note: string): ConsentRow { const existing = db .prepare("SELECT * FROM export_consents WHERE product_id = ? AND snapshot_id = ? AND snapshot_sha256 = ? AND revoked_at IS NULL ORDER BY granted_at DESC LIMIT 1") .get(productId, snapshot.id, snapshot.snapshot_sha256) as unknown as ConsentRow | undefined; if (existing) return existing; const id = newId("xpc"); db.prepare( `INSERT INTO export_consents (id, product_id, snapshot_id, snapshot_sha256, source_sha256, scope_json, note, granted_by, granted_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)` ).run(id, productId, snapshot.id, snapshot.snapshot_sha256, snapshot.source_sha256, JSON.stringify(consentScope(base)), note, actorId, nowIso()); audit(db, { actor: actorId, action: "export.consent.granted", entity: "export_consents", entityId: id, detail: { productId, snapshotSha256: snapshot.snapshot_sha256, itemCount: base.variants.length + 1, auto: true } }); return consentRow(db, productId, id); } /** Zapis wsadu: job dispatchu + pozycje z zamrożonym kontraktem. */ function insertBatch(db: Db, input: { productId: string; consent: ConsentRow; snapshot: SnapshotRow; base: PreviewBase; actorId: string }): BatchRow { const batchId = newId("xpb"); const now = nowIso(); const job = enqueueJob(db, { type: "export_dispatch", payload: { batchId, productId: input.productId }, idempotencyKey: `export:${batchId}`, maxAttempts: 5, }); db.prepare( `INSERT INTO export_batches (id, product_id, consent_id, snapshot_id, snapshot_sha256, job_id, status, created_by, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, 'queued', ?, ?, ?)` ).run(batchId, input.productId, input.consent.id, input.snapshot.id, input.snapshot.snapshot_sha256, job.id, input.actorId, now, now); const insert = db.prepare( `INSERT INTO export_items (id, batch_id, product_id, kind, variant_id, position, sku, method, payload_json, payload_sha256, status, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?, ?)` ); const addItem = (kind: "parent" | "variant", variantId: string | null, position: number, sku: string | null, params: Record) => { const payload = { method: input.base.method, parameters: params }; insert.run(newId("xpi"), batchId, input.productId, kind, variantId, position, sku, input.base.method, JSON.stringify(payload), contentHash(payload), now, now); }; addItem("parent", null, 0, String(translateParentParams(input.base).sku), translateParentParams(input.base)); input.base.variants.forEach((item, index) => addItem("variant", item.variantId, index + 1, item.sku, translateVariantParams(item, input.base))); audit(db, { actor: input.actorId, action: "export.batch.created", entity: "export_batches", entityId: batchId, detail: { productId: input.productId, consentId: input.consent.id, items: input.base.variants.length + 1, snapshotSha256: input.snapshot.snapshot_sha256 } }); return db.prepare("SELECT * FROM export_batches WHERE id = ?").get(batchId) as unknown as BatchRow; } /** Job dispatchu dla istniejącego wsadu — aktywny zostaje, martwy dostaje nowy. */ function ensureDispatchJob(db: Db, batch: BatchRow, productId: string): void { const current = batch.job_id ? getJob(db, batch.job_id) : null; if (current && ["queued", "running", "waiting"].includes(current.status)) return; const job = enqueueJob(db, { type: "export_dispatch", payload: { batchId: batch.id, productId }, idempotencyKey: `export:${batch.id}:resume:${newId("xpr")}`, maxAttempts: 5, }); db.prepare("UPDATE export_batches SET job_id = ? WHERE id = ?").run(job.id, batch.id); } function consentDto(row: ConsentRow) { return { id: row.id, snapshotId: row.snapshot_id, snapshotSha256: row.snapshot_sha256, sourceSha256: row.source_sha256, scope: JSON.parse(row.scope_json), note: row.note, grantedBy: row.granted_by, grantedAt: row.granted_at, active: row.revoked_at === null, revokedAt: row.revoked_at, revokedBy: row.revoked_by, revokeNote: row.revoke_note, }; } function itemDto(row: ItemRow, withPayload = false) { return { id: row.id, batchId: row.batch_id, kind: row.kind, variantId: row.variant_id, position: row.position, sku: row.sku, method: row.method, payloadSha256: row.payload_sha256, status: row.status, externalId: row.external_id, attempts: row.attempts, warnings: row.warnings_json ? JSON.parse(row.warnings_json) : [], lastError: row.last_error, readback: row.readback_json ? JSON.parse(row.readback_json) : null, createdAt: row.created_at, updatedAt: row.updated_at, ...(withPayload ? { payload: JSON.parse(row.payload_json) } : {}), }; } function batchDto(db: Db, row: BatchRow, withItems = true) { const items = withItems ? (db.prepare("SELECT * FROM export_items WHERE batch_id = ? ORDER BY position").all(row.id) as unknown as ItemRow[]) : []; return { id: row.id, consentId: row.consent_id, snapshotId: row.snapshot_id, snapshotSha256: row.snapshot_sha256, jobId: row.job_id, status: row.status, createdBy: row.created_by, createdAt: row.created_at, updatedAt: row.updated_at, summary: { total: items.length, verified: items.filter((i) => i.status === "verified").length, failed: items.filter((i) => i.status === "failed").length, pending: items.filter((i) => i.status === "pending").length, unconfirmed: items.filter((i) => i.status === "unconfirmed").length, cancelled: items.filter((i) => i.status === "cancelled").length, }, items: items.map((i) => itemDto(i)), }; } export function exportState(db: Db, productId: string, adapter: BaseAdapter) { const consents = db.prepare("SELECT * FROM export_consents WHERE product_id = ? ORDER BY granted_at DESC").all(productId) as unknown as ConsentRow[]; const batches = db.prepare("SELECT * FROM export_batches WHERE product_id = ? ORDER BY created_at DESC").all(productId) as unknown as BatchRow[]; return { adapter: { configured: adapter.live, name: adapter.name }, consents: consents.map(consentDto), batches: batches.map((b) => batchDto(db, b)), }; } /** * Osobna zgoda eksportowa — wiązana z dokładnym SHA zatwierdzonego * snapshotu i jego źródłami. Nie jest zgodą na publikację Allegro ani * dowodem zapisu; umożliwia utworzenie jednego wsadu outboxu. */ export function grantExportConsent( db: Db, cfg: AppConfig, input: { productId: string; actorId: string } & z.infer ) { return tx(db, () => { const requested = snapshotRow(db, input.productId, input.snapshotId); if (requested.snapshot_sha256 !== input.expectedSnapshotSha256) { throw new ExportError("SHA snapshotu nie odpowiada podglądowi — odśwież przed zgodą.", 409); } // auto-naprawa: nieaktualne źródła → nowy snapshot, brak decyzji → auto-akceptacja const { snapshot, base } = resolveExportSnapshot(db, cfg, input.productId, input.snapshotId, input.actorId); if (base.unresolved.length) { throw new ExportError(`Snapshot ma niepotwierdzone pola docelowe: ${base.unresolved.join(", ")}.`, 409); } const id = newId("xpc"); const scope = consentScope(base); db.prepare( `INSERT INTO export_consents (id, product_id, snapshot_id, snapshot_sha256, source_sha256, scope_json, note, granted_by, granted_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)` ).run(id, input.productId, snapshot.id, snapshot.snapshot_sha256, snapshot.source_sha256, JSON.stringify(scope), input.note, input.actorId, nowIso()); audit(db, { actor: input.actorId, action: "export.consent.granted", entity: "export_consents", entityId: id, detail: { productId: input.productId, snapshotSha256: snapshot.snapshot_sha256, itemCount: scope.itemCount } }); return consentDto(consentRow(db, input.productId, id)); }); } /** Odwołanie zgody kasuje niewysłane pozycje jej wsadu; wysłane zostają w historii. */ export function revokeExportConsent(db: Db, input: { productId: string; consentId: string; actorId: string; note: string }) { return tx(db, () => { const consent = consentRow(db, input.productId, input.consentId); if (consent.revoked_at) throw new ExportError("Ta zgoda została już odwołana.", 409); db.prepare("UPDATE export_consents SET revoked_at = ?, revoked_by = ?, revoke_note = ? WHERE id = ?") .run(nowIso(), input.actorId, input.note, consent.id); const batch = db.prepare("SELECT * FROM export_batches WHERE consent_id = ?").get(consent.id) as unknown as BatchRow | undefined; if (batch) { db.prepare("UPDATE export_items SET status = 'cancelled', updated_at = ? WHERE batch_id = ? AND status IN ('pending','sending','unconfirmed')") .run(nowIso(), batch.id); const left = db.prepare("SELECT COUNT(*) n FROM export_items WHERE batch_id = ? AND status NOT IN ('cancelled','verified')").get(batch.id) as { n: number }; db.prepare("UPDATE export_batches SET status = ?, updated_at = ? WHERE id = ?").run(left.n === 0 ? "cancelled" : "partial", nowIso(), batch.id); } audit(db, { actor: input.actorId, action: "export.consent.revoked", entity: "export_consents", entityId: consent.id, detail: { productId: input.productId, batchId: batch?.id ?? null } }); return consentDto(consentRow(db, input.productId, consent.id)); }); } /** * Wsad outboxu ze zgody — jeden wsad na zgodę. Pozycje zamrażają * przetłumaczony kontrakt (media jako referencje SHA, bez bajtów). */ export function createExportBatch( db: Db, cfg: AppConfig, input: { productId: string; consentId: string; actorId: string } ) { return tx(db, () => { const consent = consentRow(db, input.productId, input.consentId); if (consent.revoked_at) throw new ExportError("Zgoda eksportowa została odwołana — wymagana nowa zgoda.", 409); const existing = db.prepare("SELECT * FROM export_batches WHERE consent_id = ?").get(consent.id) as unknown as BatchRow | undefined; if (existing) return { batch: batchDto(db, existing), created: false }; // auto-naprawa: nieaktualne źródła → nowy snapshot + przebindowanie zgody, // brak akceptacji → auto-decyzja. Zamiast odrzucać żądanie. const { snapshot, base } = resolveExportSnapshot(db, cfg, input.productId, consent.snapshot_id, input.actorId); if (snapshot.id !== consent.snapshot_id || snapshot.snapshot_sha256 !== consent.snapshot_sha256 || snapshot.source_sha256 !== consent.source_sha256) { db.prepare("UPDATE export_consents SET snapshot_id = ?, snapshot_sha256 = ?, source_sha256 = ?, scope_json = ? WHERE id = ?") .run(snapshot.id, snapshot.snapshot_sha256, snapshot.source_sha256, JSON.stringify(consentScope(base)), consent.id); } const row = insertBatch(db, { productId: input.productId, consent, snapshot, base, actorId: input.actorId }); return { batch: batchDto(db, row), created: true }; }); } /** * WYSYŁKA BEZPOŚREDNIA — jedno kliknięcie z panelu produktu. * Buduje bieżącą paczkę, w locie zapisuje snapshot + auto-akceptację + * zgodę i od razu kolejkuje wsad do dispatchu. Żadnych ręcznych zgód, * wpisywania SHA ani ocen po drodze. * * Ponowne kliknięcie tego samego stanu: wznawia niedokończone pozycje * (verified nigdy nie jest wysyłane ponownie — brak duplikatów w Base), * unconfirmed pojedna runner po SKU. Gdy źródła się zmieniły, powstaje * nowy snapshot + nowa zgoda + nowy wsad; stare pozostają w historii. */ export function sendExportNow( db: Db, cfg: AppConfig, input: { productId: string; actorId: string } ) { return tx(db, () => { const current = buildOfferPackage(db, cfg, input.productId); const base = previewBaseSchema.parse(current.base); if (!base.variants.length) { throw new ExportError("Produkt nie ma wariantów do wysłania — powiąż tkaninę i strony (kroki 1 i 3).", 409); } // Kontrakt walidowany przed zapisem czegokolwiek — jeden komunikat // zbiera wszystkie braki (np. konfigurację base.* z Ustawień). const problems = new Set(); try { translateParentParams(base); } catch (e) { problems.add(e instanceof Error ? e.message : String(e)); } for (const item of base.variants) { try { translateVariantParams(item, base); } catch (e) { problems.add(`${item.side}/${item.shade}: ${e instanceof Error ? e.message : String(e)}`); } } if (problems.size) { throw new ExportError(`Do wysyłki brakuje danych: ${[...problems].join(" · ")}`, 409); } const snapshot = findOrCreateSnapshot(db, input.productId, current, input.actorId); ensureSnapshotApproved(db, snapshot, input.actorId); const consent = findOrCreateConsent(db, input.productId, snapshot, base, input.actorId, "Wysyłka bezpośrednia z panelu produktu — jeden przycisk."); const existing = db.prepare("SELECT * FROM export_batches WHERE consent_id = ?").get(consent.id) as unknown as BatchRow | undefined; if (existing) { const open = (db.prepare("SELECT COUNT(*) n FROM export_items WHERE batch_id = ? AND status != 'verified'").get(existing.id) as { n: number }).n; if (open === 0) { return { batch: batchDto(db, existing), created: false, alreadySent: true }; } const now = nowIso(); db.prepare("UPDATE export_items SET status = 'pending', last_error = NULL, updated_at = ? WHERE batch_id = ? AND status IN ('failed','cancelled')").run(now, existing.id); db.prepare("UPDATE export_batches SET status = 'queued', updated_at = ? WHERE id = ?").run(now, existing.id); ensureDispatchJob(db, existing, input.productId); audit(db, { actor: input.actorId, action: "export.batch.resumed", entity: "export_batches", entityId: existing.id, detail: { productId: input.productId, open } }); const row = db.prepare("SELECT * FROM export_batches WHERE id = ?").get(existing.id) as unknown as BatchRow; return { batch: batchDto(db, row), created: false, resumed: true }; } const row = insertBatch(db, { productId: input.productId, consent, snapshot, base, actorId: input.actorId }); return { batch: batchDto(db, row), created: true }; }); } export function getExportItem(db: Db, productId: string, itemId: string) { const row = db.prepare("SELECT * FROM export_items WHERE id = ? AND product_id = ?").get(itemId, productId) as unknown as ItemRow | undefined; if (!row) throw new ExportError("Nie znaleziono pozycji eksportu dla tego produktu.", 404); return itemDto(row, true); } /** Ponowienie pozycji failed/cancelled — wymaga aktywnej zgody; nowy job wsadu. */ export function retryExportItem(db: Db, input: { productId: string; itemId: string; actorId: string; note: string }) { return tx(db, () => { const item = db.prepare("SELECT * FROM export_items WHERE id = ? AND product_id = ?").get(input.itemId, input.productId) as unknown as ItemRow | undefined; if (!item) throw new ExportError("Nie znaleziono pozycji eksportu dla tego produktu.", 404); if (!["failed", "cancelled"].includes(item.status)) throw new ExportError("Ponowić można tylko pozycję zakończoną błędem lub anulowaną.", 409); const batch = db.prepare("SELECT * FROM export_batches WHERE id = ?").get(item.batch_id) as unknown as BatchRow; const consent = consentRow(db, input.productId, batch.consent_id); if (consent.revoked_at) throw new ExportError("Zgoda tego wsadu została odwołana — pozycja nie może wrócić do kolejki.", 409); const now = nowIso(); db.prepare("UPDATE export_items SET status = 'pending', last_error = NULL, updated_at = ? WHERE id = ?").run(now, item.id); db.prepare("UPDATE export_batches SET status = 'queued', updated_at = ? WHERE id = ?").run(now, batch.id); const job = enqueueJob(db, { type: "export_dispatch", payload: { batchId: batch.id, productId: input.productId }, idempotencyKey: `export:${batch.id}:retry:${item.id}:${item.attempts}`, maxAttempts: 5, }); db.prepare("UPDATE export_batches SET job_id = ? WHERE id = ?").run(job.id, batch.id); audit(db, { actor: input.actorId, action: "export.item.retry", entity: "export_items", entityId: item.id, detail: { productId: input.productId, batchId: batch.id, note: input.note } }); return itemDto(db.prepare("SELECT * FROM export_items WHERE id = ?").get(item.id) as unknown as ItemRow); }); } function setItem(db: Db, id: string, fields: Record) { const sets = Object.keys(fields).map((k) => `${k} = ?`).join(", "); db.prepare(`UPDATE export_items SET ${sets}, updated_at = ? WHERE id = ?`).run(...Object.values(fields), nowIso(), id); } function loadImageBytes(db: Db, cfg: AppConfig, ref: { renditionId: string; sha256: string }) { const row = getRendition(db, ref.renditionId); const file = row ? renditionFile(cfg, row) : null; if (!row || !file || row.sha256 !== ref.sha256) { throw new ExportError("Rendition pozycji nie istnieje albo nie odpowiada SHA ze snapshotu.", 409); } return Promise.resolve({ data: fs.readFileSync(file.abs).toString("base64"), bytes: row.bytes }); } function loadVideoBytes(db: Db, cfg: AppConfig, ref: { assetId: string; sha256: string }) { const asset = getAsset(db, ref.assetId); if (!asset || asset.sha256 !== ref.sha256) throw new ExportError("Klip pozycji nie odpowiada SHA ze snapshotu.", 409); const abs = path.resolve(cfg.mediaDir, asset.path); if (!abs.startsWith(path.resolve(cfg.mediaDir) + path.sep) || !fs.existsSync(abs)) { throw new ExportError("Plik klipu nie istnieje na dysku.", 409); } return Promise.resolve({ data: fs.readFileSync(abs).toString("base64"), bytes: asset.bytes }); } /** Read-back: zapis nie jest potwierdzony bez odczytu zwrotnego tego samego rekordu. */ function verifyReadBack(item: ItemRow, params: Record, rows: Array>): { ok: boolean; detail: Record } { const found = rows.find((r) => Number(r.id ?? r.product_id) === item.external_id); if (!found) return { ok: false, detail: { reason: "Brak product_id w odczycie zwrotnym." } }; const remoteSku = typeof found.sku === "string" ? found.sku : null; if (params.sku && remoteSku !== params.sku) { return { ok: false, detail: { reason: `SKU odczytu (${remoteSku ?? "brak"}) ≠ wysłane (${String(params.sku)}).` } }; } const textFields = (found.text_fields && typeof found.text_fields === "object" ? found.text_fields : {}) as Record; const sentName = (params.text_fields as Record | undefined)?.name; if (sentName && typeof textFields.name === "string" && textFields.name !== sentName) { return { ok: false, detail: { reason: "Nazwa w odczycie zwrotnym nie odpowiada wysłanej." } }; } // EAN: jeśli wariant wysłał kod, rekord zdalny musi go zwrócić — // brakujący/cudzy EAN to cicha korupcja identyfikatora w magazynie. const sentEan = typeof params.ean === "string" && params.ean ? params.ean : null; if (sentEan) { const remoteEan = typeof found.ean === "string" && found.ean ? found.ean : null; if (remoteEan !== sentEan) { return { ok: false, detail: { reason: `EAN odczytu (${remoteEan ?? "brak"}) ≠ wysłane (${sentEan}).` } }; } } // Zdjęcia: wysłana galeria musi być widoczna w odczycie — sloty // pomijane przy limicie/błędzie dekodowania nie mogą dawać „verified". const sentImageCount = Array.isArray(params.images) ? params.images.length : 0; if (sentImageCount > 0) { const remote = found.images; const remoteCount = Array.isArray(remote) ? remote.length : remote && typeof remote === "object" ? Object.keys(remote as Record).length : 0; if (remoteCount === 0) { return { ok: false, detail: { reason: `Odczyt zwrotny nie zawiera zdjęć, mimo że wysłano ${sentImageCount} pozycji galerii.` } }; } } return { ok: true, detail: { productId: item.external_id, sku: remoteSku, name: textFields.name ?? null, ean: found.ean ?? null, images: sentImageCount } }; } /** * Wykonawca wsadu outboxu. Kolejność: rodzic → warianty (parent_id z * zweryfikowanego external_id). Każda pozycja ma własny status; awaria * jednej nie cofa pozostałych. Timeout po wysłaniu → 'unconfirmed' i * pojednanie po SKU przy następnym przebiegu — nigdy ślepy re-create. */ export async function runExportBatch( db: Db, cfg: AppConfig, adapter: BaseAdapter, jobId: string, isStopped: () => boolean = () => false ): Promise { const job = getJob(db, jobId); if (!job) return; const payload = z.object({ batchId: z.string(), productId: z.string() }).strict().parse(JSON.parse(job.payload_json)); const batch = db.prepare("SELECT * FROM export_batches WHERE id = ? AND product_id = ?").get(payload.batchId, payload.productId) as unknown as BatchRow | undefined; if (!batch || batch.job_id !== jobId) { finishJob(db, jobId, { ok: false, error: "Zadanie nie jest przypisane do aktywnego wsadu eksportu." }); return; } if (["done", "cancelled"].includes(batch.status)) { finishJob(db, jobId, { ok: true, result: { batchId: batch.id, status: batch.status } }); return; } // Awaria w trakcie wysyłki: wynik nieznany → pojednanie, nie re-create. db.prepare("UPDATE export_items SET status = 'unconfirmed', updated_at = ? WHERE batch_id = ? AND status = 'sending'").run(nowIso(), batch.id); // Ponowny dispatch sam naprawia wsad: pozycje zakończone błędem wracają // do kolejki, zamiast wiecznie zostawać 'failed' przy pominiętych // następnych pozycjach. 'verified' i 'cancelled' są nietknięte — // verified nigdy nie leci ponownie, cancelled wymaga jawnej decyzji // (sendExportNow / retryExportItem), a odwołaną zgodę i tak domyka gate. db.prepare("UPDATE export_items SET status = 'pending', last_error = NULL, updated_at = ? WHERE batch_id = ? AND status = 'failed'").run(nowIso(), batch.id); db.prepare("UPDATE export_batches SET status = 'running', updated_at = ? WHERE id = ?").run(nowIso(), batch.id); const consent = db.prepare("SELECT * FROM export_consents WHERE id = ?").get(batch.consent_id) as unknown as ConsentRow; const inventoryId = (JSON.parse(consent.scope_json) as { inventoryId: number }).inventoryId; const items = db.prepare("SELECT * FROM export_items WHERE batch_id = ? ORDER BY position").all(batch.id) as unknown as ItemRow[]; const fail = (item: ItemRow, message: string) => setItem(db, item.id, { status: "failed", last_error: message, attempts: item.attempts + 1 }); let parentExternalId = items.find((i) => i.kind === "parent" && i.status === "verified")?.external_id ?? null; let uncertain = false; // Reweryfikacja przed każdym zapisem zewnętrznym: zgoda nadal aktywna, // snapshot nadal zatwierdzony dla tego SHA, pozycja nadal wysyłalna // (równoległe odwołanie/retry mogły zmienić stan między iteracjami). const stillAuthorized = (item: ItemRow): { ok: boolean; skip: boolean; reason?: string } => { const freshConsent = db.prepare("SELECT revoked_at FROM export_consents WHERE id = ?").get(batch.consent_id) as { revoked_at: string | null } | undefined; if (!freshConsent || freshConsent.revoked_at) return { ok: false, skip: true, reason: "Zgoda eksportowa została odwołana w trakcie wsadu." }; const decision = latestDecision(db, batch.snapshot_id); if (decision?.decision !== "approved" || decision.snapshot_sha256 !== batch.snapshot_sha256) { return { ok: false, skip: true, reason: "Snapshot paczki stracił akceptację dla SHA wsadu." }; } const freshItem = db.prepare("SELECT status FROM export_items WHERE id = ?").get(item.id) as { status: string } | undefined; if (!freshItem || !["pending", "unconfirmed"].includes(freshItem.status)) return { ok: false, skip: true }; return { ok: true, skip: false }; }; for (const item of items) { if (isStopped()) break; if (["verified", "failed", "cancelled"].includes(item.status)) continue; const gate = stillAuthorized(item); if (!gate.ok) { if (gate.skip && gate.reason) setItem(db, item.id, { status: "cancelled", last_error: gate.reason }); continue; } try { if (!adapter.live) throw new ExportBlocked("Adapter Base nie jest skonfigurowany (brak BASELINKER_API_TOKEN). Wsad pozostaje w outboxie."); const parsed = JSON.parse(item.payload_json) as { method: string; parameters: Record }; const params = { ...parsed.parameters }; if (item.status === "unconfirmed") { const found = item.sku ? await adapter.findProductIdBySku(inventoryId, item.sku) : null; if (found === null) { setItem(db, item.id, { status: "pending" }); } else { setItem(db, item.id, { status: "sending", external_id: found }); } item.status = found === null ? "pending" : "sending"; if (found !== null) item.external_id = found; } if (item.status === "pending") { if (item.kind === "variant") { if (parentExternalId === null) throw new ExportError("Rodzic wariantów nie jest zweryfikowany — wariant czeka na jego product_id.", 409); params.parent_id = parentExternalId; } // Idempotencja po SKU: produkt już istniejący w magazynie nie jest // duplikowany — dostaje update (addInventoryProduct z product_id // to edycja rekordu), więc treść/zdjęcia odnawiają się na żywo. if (item.sku && item.external_id === null) { const found = await adapter.findProductIdBySku(inventoryId, item.sku); if (found !== null) { setItem(db, item.id, { external_id: found }); item.external_id = found; } } { // Niekompletne referencje zdjęć (rendition niegotowy) nie blokują // wysyłki — pozycja leci z tym, co jest; pominięcia w warnings. const skipped: string[] = []; if (Array.isArray(params.images)) { const refs = params.images as Array<{ position: number; renditionId: string | null; sha256: string | null }>; const ready = refs.filter((r) => r.renditionId && r.sha256); const dropped = refs.length - ready.length; if (dropped > 0) skipped.push(`Pominięto ${dropped} niegotowych zdjęć (brak renditionu w paczce).`); if (ready.length) params.images = ready; else delete params.images; } const wire = await materializeMedia(params, (ref) => loadImageBytes(db, cfg, ref), (ref) => loadVideoBytes(db, cfg, ref)); if (item.external_id !== null) wire.product_id = item.external_id; setItem(db, item.id, { status: "sending", attempts: item.attempts + 1 }); item.attempts += 1; const sent = await adapter.addInventoryProduct(wire); setItem(db, item.id, { external_id: sent.productId, warnings_json: JSON.stringify([...sent.warnings, ...skipped]) }); item.external_id = sent.productId; } } if (item.external_id === null) throw new ExportUncertain("Brak product_id po wysłaniu — pojednanie przy następnym przebiegu."); const readBack = await adapter.getInventoryProductsData(inventoryId, [item.external_id]); const check = verifyReadBack({ ...item }, params, readBack); if (!check.ok) throw new ExportRejected(`Read-back nie potwierdza zapisu: ${check.detail.reason}`); // last_error=NULL — pozycja pojednana po wcześniejszym „unconfirmed" // nie może nosić nieaktualnego komunikatu o rozbieżności. setItem(db, item.id, { status: "verified", last_error: null, readback_json: JSON.stringify({ ...check.detail, checkedAt: nowIso() }) }); if (item.kind === "parent") parentExternalId = item.external_id; } catch (err) { if (err instanceof ExportUncertain) { setItem(db, item.id, { status: "unconfirmed", last_error: err.message }); uncertain = true; break; } const message = err instanceof Error ? err.message : String(err); if (err instanceof ExportRejected) { // Zapis fizycznie doszedł (mamy product_id), ale odczyt zwrotny // się nie zgadza (SKU/nazwa/EAN/galeria). To NIE jest „verified" // ani zwykły „failed" — pozycja idzie w 'unconfirmed': kolejny // dispatch pojedna ją po SKU i zweryfikuje ponownie, a operator // widzi czytelny powód rozbieżności zamiast cichego sukcesu. setItem(db, item.id, { status: "unconfirmed", last_error: message }); } else { fail(item, message); } if (item.kind === "parent" || err instanceof ExportBlocked) break; } } if (isStopped()) { deferJob(db, jobId, 0); return; } const done = db.prepare("SELECT * FROM export_items WHERE batch_id = ?").all(batch.id) as unknown as ItemRow[]; const counts = { verified: done.filter((i) => i.status === "verified").length, failed: done.filter((i) => i.status === "failed").length, pending: done.filter((i) => i.status === "pending").length, unconfirmed: done.filter((i) => i.status === "unconfirmed").length, cancelled: done.filter((i) => i.status === "cancelled").length, }; const status = counts.cancelled === done.length ? "cancelled" : counts.verified === done.length ? "done" : counts.failed === done.length ? "failed" : "partial"; db.prepare("UPDATE export_batches SET status = ?, updated_at = ? WHERE id = ?").run(status, nowIso(), batch.id); audit(db, { actor: batch.created_by, action: "export.batch.finished", entity: "export_batches", entityId: batch.id, detail: { productId: batch.product_id, status, ...counts } }); if (counts.failed || counts.unconfirmed || counts.pending) { finishJob(db, jobId, { ok: false, error: `Wsad nie jest kompletny: ${counts.verified} zweryfikowane, ${counts.failed} błędów, ${counts.unconfirmed} niepotwierdzonych, ${counts.pending} oczekujących.`, retryable: uncertain, }); } else { finishJob(db, jobId, { ok: true, result: { batchId: batch.id, ...counts } }); } }