test bugs fixed

This commit is contained in:
2026-08-28 18:18:50 +05:30
parent 670d193f36
commit 32c3bef3ad
6 changed files with 598 additions and 473 deletions

View File

@@ -1,125 +1,149 @@
/**
* The catalogue ingest service — `mcp.nearle.ai.in`, pipeline v3.2.0.
* The catalogue ingest service — `mcp.nearle.ai.in`.
*
* A workbook goes up, the service parses and enriches it, and the rows land in
* the global catalogue's per-brand tables. It replaces the row-by-row create
* loop the console used to run in the browser.
* A workbook goes up, the service runs its eleven-stage pipeline over every
* row, and the products land in the global catalogue's per-brand tables. From
* there Fiesta already sees them: `/web/catalogue/getbrands` and
* `/web/catalogue/getproducts` read the SAME database this service writes to,
* which is why an upload here shows up in the console's catalogue with nothing
* in between to build or synchronise.
*
* Everything here is written against the contract the owning team supplied, not
* against a guess. Where a decision looks arbitrary it usually is not — the
* reason is in the comment.
* ── Batches, not jobs ────────────────────────────────────────────────────────
*
* ── Submit, then poll ────────────────────────────────────────────────────────
* This module used to call `/api/admin/store-catalog/*` — one file, one
* `job_id`, jobs held in memory and lost on restart. That API is still
* deployed, but the owning team's guidance is to integrate against
* `/api/uploads/catalog`, which takes up to twenty files at once and answers
* with a `batch_id` that survives a restart.
*
* `POST /ingest` answers 202 with a job id and hands off to a background
* thread; the outcome arrives from `GET /jobs/{id}`. Two operational details
* from the owning team shape `pollJob` below, and both are the kind of thing
* that silently produces a wrong screen if ignored:
* `POST` answers 202 the moment the files are staged; the outcome arrives from
* `GET /api/uploads/catalog/{batch_id}`. Two things about that shape earn their
* own handling below:
*
* - **Jobs live in memory.** A backend restart loses them and polling returns
* 404. That is "we no longer know", NOT "it failed" — the rows may well have
* been written. Reporting a failure there would send someone re-uploading a
* sheet that already landed.
* - **`products_built > 0` is not success.** A job whose rows were built but
* could not be stored is marked `failed` with the reason in
* `result.storage_error`. The counts are populated either way, so reading
* them without checking `status` reports an import that never happened.
* - **A file that could not be read stays in the batch** as a failed member
* rather than being dropped, so a sender who submitted two files and sees
* one knows what became of the other.
* - **`partial` is not `done`.** Some files landed and some did not, and four
* of five succeeding must never render as flat success.
*
* ── The credential never reaches this file ───────────────────────────────────
*
* `X-API-Key`, attached by the Vite proxy from `INGEST_TOKEN` in `.env.local`.
* The key carries `require_admin`, which on that service is a superuser — the
* same key reaches `/api/catalog/generate`, `/api/system/init` and the ML
* training endpoints. A key in the browser bundle is a key published to every
* visitor, so it stays server-side and this module never sees one.
* `X-API-Key`, attached by the Vite proxy in development and by nginx in
* production, both from `INGEST_TOKEN`. It stays server-side because the bundle
* is served to anyone who opens the console.
*
* That is why production is not solved here. The console's origin is not in
* their `API_CORS_ORIGINS` and should not be added: the owning team's own
* recommendation is server-to-server, which means Fiesta relays the call. Until
* that exists, this path works in development only.
* `INGEST_TOKEN` is the SECRET ALONE — the 43-character value. The service
* holds `name:role:secret` triples in its own `API_KEYS` and looks keys up by
* the secret, so pasting the whole triple fails as "Invalid API key." rather
* than as something that names the real mistake.
*/
const INGEST_BASE = import.meta.env['VITE_INGEST_BASE'] ?? '/ingest';
const ROOT = '/api/admin/store-catalog';
const ROOT = '/api/uploads/catalog';
/* ── Limits, mirroring the service's own ──────────────────────────────────── */
/**
* Client-side limits, mirroring the service's own.
*
* Checked here so a 12 MB workbook is refused in the browser instead of being
* uploaded over a shop's connection to earn a 413. The service remains the
* authority; this is politeness, not validation.
* Checked here so a drop that cannot possibly be accepted is refused in the
* browser rather than uploaded over a shop's connection to earn a 413. The
* service remains the authority; this is politeness, not validation.
*/
export const MAX_FILES = 20;
export const MAX_FILE_BYTES = 10 * 1024 * 1024;
export const MAX_TOTAL_BYTES = 50 * 1024 * 1024;
/** Per file. A sheet over this is marked failed; the rest of the batch runs. */
export const MAX_ROWS = 2000;
/** Everything the service parses. `.txt`/`.tab` included because it takes them. */
export const ACCEPTED_EXTENSIONS = ['.xlsx', '.xlsm', '.xls', '.csv', '.tsv', '.txt', '.tab'];
/** Everything the service parses, from the documented format list. */
export const ACCEPTED_EXTENSIONS = ['.xlsx', '.xls', '.csv', '.tsv'];
/* ── Response types, from the owning team's real output ───────────────────── */
/* ── Response types, from the owning team's documented output ─────────────── */
export type JobStatus = 'pending' | 'running' | 'done' | 'failed';
/**
* Seven states, not four.
*
* `partial` and `interrupted` are the two that matter and the two a client is
* most likely to collapse into something else. `interrupted` means a restart
* cut the batch short; it never auto-restarts and needs an admin to resume it,
* so reporting it as `failed` would send someone re-uploading a batch that is
* waiting to be continued.
*/
export type BatchStatus =
| 'queued'
| 'running'
| 'done'
| 'partial'
| 'failed'
| 'interrupted'
| 'cancelled';
/** A row that could not be turned into a product at all. */
export interface IngestRowError {
/** 1-based spreadsheet row. The header is row 1, so this matches Excel. */
row: number;
product_name: string;
error: string;
}
export type BatchFileStatus = 'queued' | 'running' | 'done' | 'failed';
/** A row that was built, then failed the deterministic validation gate. */
export interface IngestRejection {
product_name: string;
size: string;
reason: string;
}
export interface IngestJobResult {
rows_total: number;
/** What the pipeline made of one file, once it has finished. */
export interface BatchFileResult {
rows_total?: number;
/** Can exceed `rows_total`: "100g, 200g, 500g" in one cell is three products. */
products_built?: number;
inserted?: number;
/** Existing rows whose blank columns this upload filled in. */
backfilled?: number;
/** Already present and already complete — nothing to do. */
skipped_existing?: number;
rejected?: number;
/** Sheet header → the field it was read as. */
recognised_columns?: Record<string, string>;
/** Headers that matched nothing. Reported, never an error. */
unrecognised_columns?: string[];
/** Non-null means rows were built but never stored. */
storage_error?: string | null;
}
export interface BatchFile {
index: number;
filename: string;
status: BatchFileStatus;
/** Present on a file the service refused to read, and the reason it gives. */
detail?: string | null;
size_bytes?: number;
rows_total?: number;
/** Progress through the eleven stages, while it runs. */
stage_index?: number;
stage_name?: string;
total_stages?: number;
rows_done?: number;
result?: BatchFileResult | null;
}
export interface BatchTotals {
rows_total: number;
products_built: number;
inserted: number;
/** Existing rows whose blank columns this upload filled in. */
backfilled: number;
/** Already present and already complete — nothing to do. */
skipped_existing: number;
rejected: number;
brands: string[];
/** Capped at 50 by the service; `rejected` carries the true total. */
rejections: IngestRejection[];
/** Capped at 50 by the service; `error_count` carries the true total. */
errors: IngestRowError[];
error_count: number;
/** Non-fatal corrections that were applied anyway, row-numbered. */
warnings: string[];
/** Sheet header → the field it was read as. */
recognised_columns: Record<string, string>;
/** Headers that matched nothing and were silently dropped. */
unrecognised_columns: string[];
/** Non-null means the rows were built but never stored. */
storage_error: string | null;
}
export interface IngestJob {
job_id: string;
filename: string;
status: JobStatus;
export interface IngestBatch {
batch_id: string;
status: BatchStatus;
detail: string | null;
stage_index: number;
stage_name: string;
total_stages: number;
rows_done: number;
rows_total: number;
result: IngestJobResult | null;
}
/** What `/preview` answers — a dry run that writes nothing. */
export interface IngestPreview {
recognised_columns: Record<string, string>;
unrecognised_columns: string[];
rows: Record<string, unknown>[];
rows_total?: number;
submitted_by?: string;
/** Epoch SECONDS, not milliseconds — multiply before handing to `Date`. */
created_at?: number;
updated_at?: number;
files_total: number;
files_done: number;
files_failed: number;
current_file?: string | null;
use_llm?: boolean;
fetch_images?: boolean;
totals: BatchTotals;
/** The brands this batch touched — the way back into the catalogue view. */
brands: string[];
files: BatchFile[];
/** Present only on the POST response; the polling reads omit it. */
message?: string;
}
export class IngestError extends Error {
@@ -134,66 +158,95 @@ export class IngestError extends Error {
}
}
/** True when a poll found no such job — see the note about in-memory jobs. */
export class JobVanishedError extends IngestError {
constructor(jobId: string) {
super(
'The service no longer knows about this job — it restarts with jobs held in memory. The products may well have been written; check the catalogue before uploading again.',
404,
);
this.name = 'JobVanishedError';
this.jobId = jobId;
}
readonly jobId: string;
}
/* ── Requests ─────────────────────────────────────────────────────────────── */
export interface SubmitOptions {
file: File;
files: File[];
/**
* Default FALSE, and that is not caution — production reports
* `"ollama": false`, so the LLM path is not available there. Asking for it
* buys nothing and the enrichment that matters (HSN, price band, FSSAI, SKU)
* is deterministic lookup rather than generation.
*/
useLlm?: boolean;
/**
* Default FALSE. Image fetching is a network round trip per row evaluating up
* to 24 candidates, and it is the stage that turns ten seconds into minutes.
* Worth turning on deliberately, not by default.
* Default FALSE. Stage 6 spawns a Playwright subprocess and searches for an
* image per row — minutes per batch on one vCPU. Worth turning on
* deliberately, never by default.
*/
fetchImages?: boolean;
/**
* Default FALSE, and not caution: production runs with `USE_OLLAMA=false`, so
* asking for it is a documented no-op. Left as an option only so the flag is
* not silently unavailable the day that changes.
*/
useLlm?: boolean;
signal?: AbortSignal;
}
function guardFile(file: File): void {
if (file.size > MAX_FILE_BYTES) {
/**
* Refuses a drop the service is certain to reject, and says which file is at
* fault rather than reporting the batch as generically too large.
*/
function guardFiles(files: File[]): void {
if (files.length === 0) {
throw new IngestError('Choose at least one file.', 400);
}
if (files.length > MAX_FILES) {
throw new IngestError(
`That file is ${(file.size / 1024 / 1024).toFixed(1)} MB. The service accepts up to 10 MB.`,
`That is ${files.length} files. The service takes ${MAX_FILES} per upload — send them in smaller batches.`,
413,
);
}
if (file.size === 0) {
throw new IngestError('That file is empty.', 400);
for (const file of files) {
if (file.size === 0) {
throw new IngestError(`"${file.name}" is empty.`, 400);
}
if (file.size > MAX_FILE_BYTES) {
throw new IngestError(
`"${file.name}" is ${(file.size / 1024 / 1024).toFixed(1)} MB. The limit is 10 MB per file.`,
413,
);
}
/**
* The FILENAME picks the parser, not the bytes. A name the service does not
* recognise comes back as an unexplained parse failure, so it is named here
* instead.
*/
const extension = file.name.toLowerCase().slice(file.name.lastIndexOf('.'));
if (!file.name.includes('.') || !ACCEPTED_EXTENSIONS.includes(extension)) {
throw new IngestError(
`"${file.name}" is not a format the service reads. Accepted: ${ACCEPTED_EXTENSIONS.join(', ')}.`,
400,
);
}
}
const total = files.reduce((sum, file) => sum + file.size, 0);
if (total > MAX_TOTAL_BYTES) {
throw new IngestError(
`That is ${(total / 1024 / 1024).toFixed(1)} MB in total. The limit is 50 MB per upload.`,
413,
);
}
}
async function send<T>(path: string, file: File, query: URLSearchParams, signal?: AbortSignal) {
guardFile(file);
/**
* Submits the sheets. Answers 202 with a batch to poll — it does not wait.
*
* The form field is `files` and it is REPEATED, once per file. The older
* endpoint took a single `file`, and sending that name here parses as no files
* at all.
*/
export async function submitBatch(options: SubmitOptions): Promise<IngestBatch> {
const { files, fetchImages = false, useLlm = false, signal } = options;
guardFiles(files);
const form = new FormData();
// `file`, confirmed: the handler signature is `file: UploadFile = File(...)`.
form.append('file', file, file.name);
for (const file of files) form.append('files', file, file.name);
// NOTE: no tenantid/locationid. This endpoint writes the GLOBAL catalogue and
// has no concept of an outlet — making a product sellable at a shop is a
// separate call to `/api/upload/stores`, keyed on `image_id`. Sending them
// here achieved nothing and implied a link that does not exist.
const query = new URLSearchParams({
fetch_images: String(fetchImages),
use_llm: String(useLlm),
});
let response: Response;
try {
response = await fetch(`${INGEST_BASE}${ROOT}${path}?${query}`, {
response = await fetch(`${INGEST_BASE}${ROOT}?${query}`, {
method: 'POST',
body: form,
// Content-Type is deliberately unset: the browser adds it WITH the
@@ -211,9 +264,54 @@ async function send<T>(path: string, file: File, query: URLSearchParams, signal?
);
}
return readResponse<T>(response);
return readResponse<IngestBatch>(response);
}
/** One poll. */
export async function fetchBatch(batchId: string, signal?: AbortSignal): Promise<IngestBatch> {
let response: Response;
try {
response = await fetch(`${INGEST_BASE}${ROOT}/${encodeURIComponent(batchId)}`, {
headers: { Accept: 'application/json' },
...(signal ? { signal } : {}),
});
} catch {
throw new IngestError('Lost contact with the ingest service while waiting.', 0);
}
return readResponse<IngestBatch>(response);
}
/** True once the batch has stopped moving, whatever the outcome. */
export function isSettled(batch: IngestBatch): boolean {
return batch.status !== 'queued' && batch.status !== 'running';
}
/**
* Polls until the batch settles.
*
* Every two seconds. The pipeline's own stages take far longer than that, and a
* person is watching a progress bar — a slower cadence buys nothing but a
* screen that looks stuck. `onTick` fires on each reading so the caller can
* render the stage name and row counts as they move.
*/
export async function pollBatch(
batchId: string,
onTick: (batch: IngestBatch) => void,
signal?: AbortSignal,
): Promise<IngestBatch> {
for (;;) {
if (signal?.aborted) throw new IngestError('Cancelled.', 0);
const batch = await fetchBatch(batchId, signal);
onTick(batch);
if (isSettled(batch)) return batch;
await new Promise((resolve) => setTimeout(resolve, 2000));
}
}
/* ── Reading a response ───────────────────────────────────────────────────── */
async function readResponse<T>(response: Response): Promise<T> {
const text = await response.text();
let payload: unknown = null;
@@ -227,6 +325,37 @@ async function readResponse<T>(response: Response): Promise<T> {
return payload as T;
}
/**
* `detail` is a STRING on some failures and an OBJECT on others.
*
* 400 and 413 send a sentence; 422 sends `{message, rows_total, errors[]}`.
* Rendering it straight prints "[object Object]" for exactly the response that
* carries the most useful information, so both shapes are unpacked here.
*/
function detailOf(payload: unknown): string | undefined {
if (payload === null || typeof payload !== 'object') return undefined;
const detail = (payload as { detail?: unknown }).detail;
if (typeof detail === 'string') return detail;
if (detail !== null && typeof detail === 'object') {
const nested = detail as { message?: unknown; errors?: unknown };
const message = typeof nested.message === 'string' ? nested.message : undefined;
const errors = Array.isArray(nested.errors) ? nested.errors : [];
// Row numbers are what makes a 422 actionable — they are the sheet's own
// 1-based numbering, header included, so they match what the operator sees.
const rows = errors
.slice(0, 5)
.map((entry) => {
const row = (entry as { row?: unknown }).row;
const error = (entry as { error?: unknown }).error;
return `row ${String(row)}: ${String(error)}`;
})
.join(' · ');
return [message, rows].filter(Boolean).join(' — ') || undefined;
}
return undefined;
}
/**
* The service's failures, in words that name the fix.
*
@@ -235,163 +364,92 @@ async function readResponse<T>(response: Response): Promise<T> {
* deployment.
*/
function describe(status: number, payload: unknown, text: string): IngestError {
const detail =
(payload !== null &&
typeof payload === 'object' &&
typeof (payload as { detail?: unknown }).detail === 'string' &&
(payload as { detail: string }).detail) ||
undefined;
const detail = detailOf(payload);
const body = text.slice(0, 2000);
if (status === 401 || status === 403) {
if (status === 401) {
// Deliberately NOT `detail ?? …`: the service answers both a missing key
// and a malformed one with a flat "Invalid API key.", which is true and
// tells nobody what to change.
return new IngestError(
'The ingest service rejected the credential. INGEST_TOKEN must be the SECRET ONLY — the 43-character value, not the `name:role:secret` triple, which fails as an invalid key rather than a malformed one. Set it on the container and restart; envsubst runs at container start, so a running container will not pick it up.',
status,
body,
);
}
if (status === 403) {
return new IngestError(
detail ??
'The ingest service rejected the credential. Set INGEST_TOKEN in .env.local and restart the dev server — and note the key only works once their backend is rebuilt with it, since API_KEYS is baked in at build time.',
'That key authenticated but does not hold `upload_catalog`. It needs the `uploader` or `admin` role.',
status,
text.slice(0, 2000),
body,
);
}
if (status === 413) {
return new IngestError(
detail ?? 'Too large for the service — the limits are 10 MB and 2000 rows.',
detail ?? 'Too large for the service — 20 files, 10 MB each, 50 MB and 20,000 rows per upload.',
status,
text.slice(0, 2000),
body,
);
}
if (status === 429) {
return new IngestError(
detail ??
'Four batches are already queued. This upload was staged rather than lost — wait a moment and send it again.',
status,
body,
);
}
if (status === 400) {
return new IngestError(
detail ?? 'The service could not read that file — it may be empty or have no data rows.',
detail ??
'The service could not read that file. Every sheet needs a product-name column — product, item, variant or name.',
status,
text.slice(0, 2000),
body,
);
}
return new IngestError(detail ?? `The ingest service returned HTTP ${status}.`, status, text.slice(0, 2000));
return new IngestError(detail ?? `The ingest service returned HTTP ${status}.`, status, body);
}
/**
* A true dry run. Parses, reports the column mapping and the first rows, and
* writes nothing at all.
*
* Run before every ingest. It is the only way to see `unrecognised_columns`
* before the fact, and a header that matched nothing is dropped SILENTLY — a
* price column the service never saw looks exactly like a successful import
* until someone opens the catalogue.
*/
export function previewSheet(file: File, signal?: AbortSignal): Promise<IngestPreview> {
return send<IngestPreview>('/preview', file, new URLSearchParams(), signal);
}
/* ── Reading a finished batch ─────────────────────────────────────────────── */
/** Submits the sheet. Answers 202 with a job to poll — it does not wait. */
export function submitIngest(options: SubmitOptions): Promise<IngestJob> {
const { file, useLlm = false, fetchImages = false, signal } = options;
const query = new URLSearchParams({
use_llm: String(useLlm),
fetch_images: String(fetchImages),
});
return send<IngestJob>('/ingest', file, query, signal);
}
/** One poll. Throws `JobVanishedError` on 404 — see the note at the top. */
export async function fetchJob(jobId: string, signal?: AbortSignal): Promise<IngestJob> {
let response: Response;
try {
response = await fetch(`${INGEST_BASE}${ROOT}/jobs/${encodeURIComponent(jobId)}`, {
headers: { Accept: 'application/json' },
...(signal ? { signal } : {}),
});
} catch {
throw new IngestError('Lost contact with the ingest service while waiting.', 0);
}
if (response.status === 404) throw new JobVanishedError(jobId);
return readResponse<IngestJob>(response);
}
/**
* Polls until the job settles.
*
* Every second. The service does no rate limiting on this and a person is
* watching a progress bar, so a slower cadence buys nothing but a screen that
* looks stuck. `onTick` fires on each reading so the caller can render
* `stage_name` and `rows_done` as they move.
*/
export async function pollJob(
jobId: string,
onTick: (job: IngestJob) => void,
signal?: AbortSignal,
): Promise<IngestJob> {
for (;;) {
if (signal?.aborted) throw new IngestError('Cancelled.', 0);
const job = await fetchJob(jobId, signal);
onTick(job);
if (job.status === 'done' || job.status === 'failed') return job;
await new Promise((resolve) => setTimeout(resolve, 1000));
}
}
/* ── Deriving what the service does not return ────────────────────────────── */
/**
* The catalogue's primary key, computed locally.
*
* The ingest returns counts, not ids — but the key is deterministic, so the
* rows can be addressed without being told. That matters for the step after
* this one: `/api/upload/stores` joins on exactly this value to put a product
* on a shop's shelf.
*
* image_id = sanitize(brand) + "_" + slugify(name [+ " " + size])
*
* The size is appended ONLY when its slug is not already inside the name's —
* "Good Day Cashew Cookies 100g" with size "100g" must not become
* `..._100g_100g`.
*
* Verified against real output: `britannia_britannia_good_day_cashew_cookies_100g`.
*/
export function imageId(brand: string, productName: string, size?: string): string {
const nameSlug = slugify(productName);
const sizeSlug = size ? slugify(size) : '';
const tail = sizeSlug && !nameSlug.includes(sizeSlug) ? `${nameSlug}_${sizeSlug}` : nameSlug;
return `${sanitize(brand)}_${tail}`;
}
/** lowercase · space, hyphen and & become `_` · drop the rest · collapse runs. */
function sanitize(value: string): string {
return value
.toLowerCase()
.replace(/[\s\-&]/g, '_')
.replace(/[^a-z0-9_]/g, '')
.replace(/_{2,}/g, '_')
.replace(/^_+|_+$/g, '');
}
/** lowercase · any run of non-alphanumerics becomes one `_` · trim. */
function slugify(value: string): string {
return value
.toLowerCase()
.replace(/[^a-z0-9]+/g, '_')
.replace(/^_+|_+$/g, '');
}
/* ── Reading a finished job ───────────────────────────────────────────────── */
/** True when the job ended without the rows reaching the database. */
export function isStorageFailure(job: IngestJob): boolean {
return job.status === 'failed' || Boolean(job.result?.storage_error);
/** True when the batch ended without everything landing. */
export function isIncomplete(batch: IngestBatch): boolean {
return (
batch.status === 'partial' ||
batch.status === 'failed' ||
batch.status === 'interrupted' ||
batch.status === 'cancelled' ||
batch.files_failed > 0
);
}
/** One line for the top of the result panel. */
export function summarise(job: IngestJob): string {
const result = job.result;
if (!result) return job.detail ?? 'The service returned no result.';
export function summarise(batch: IngestBatch): string {
const { totals } = batch;
if (isStorageFailure(job)) {
return result.storage_error
? `Built ${result.products_built} products but could not store them: ${result.storage_error}`
: (job.detail ?? 'The job failed.');
if (batch.status === 'failed') {
return batch.detail ?? 'No file could be ingested.';
}
if (batch.status === 'interrupted') {
return 'The service restarted part-way through. An admin can resume this batch — it will not restart on its own.';
}
if (batch.status === 'cancelled') {
return 'This batch was cancelled before every file ran.';
}
const parts = [`${result.inserted} added`];
if (result.backfilled > 0) parts.push(`${result.backfilled} filled in`);
if (result.skipped_existing > 0) parts.push(`${result.skipped_existing} already there`);
return parts.join(' · ');
const parts = [`${totals?.inserted ?? 0} added`];
if ((totals?.backfilled ?? 0) > 0) parts.push(`${totals.backfilled} filled in`);
if ((totals?.skipped_existing ?? 0) > 0) parts.push(`${totals.skipped_existing} already there`);
if ((totals?.rejected ?? 0) > 0) parts.push(`${totals.rejected} rejected`);
const summary = parts.join(' · ');
return batch.files_failed > 0
? `${summary} — but ${batch.files_failed} of ${batch.files_total} files could not be read`
: summary;
}
/** Overall progress, for a bar. Stages within a file are too fine to show. */
export function progressOf(batch: IngestBatch): { done: number; total: number } {
return { done: batch.files_done + batch.files_failed, total: batch.files_total };
}