Files
daily_console_web/src/api/ingest.ts

677 lines
26 KiB
TypeScript

/**
* The catalogue ingest service — `mcp.nearle.ai.in`.
*
* A spreadsheet goes up, an admin reviews it, and once released the eleven-stage
* pipeline writes the products into the global catalogue. From there Fiesta
* already sees them: `/web/catalogue/getbrands` and `/web/catalogue/getproducts`
* read the SAME database the pipeline writes to, so an upload appears in the
* console with nothing in between to build or synchronise.
*
* ── A drop is not a run ──────────────────────────────────────────────────────
*
* `POST /api/uploads/catalog` creates a DROP, and nothing runs on arrival. The
* files wait in an admin review inbox; only when someone selects them and
* presses Start does a RUN begin, under a different id. The drop id stays valid
* for the whole lifecycle and its per-file status is how you follow it:
*
* queued — still in the inbox, nobody has looked
* released — accepted; `released_to` is the run, and the results are there
* dismissed — declined; nothing further is coming
*
* `resolveBatch` below makes that hop automatically, so callers poll one id and
* get whichever record actually has the answer.
*
* ── No credential ────────────────────────────────────────────────────────────
*
* The drop endpoint takes none, and that is safe precisely because of the review
* gate: an unwanted drop costs disk until somebody declines it, never products
* in the live catalogue.
*
* So `INGEST_TOKEN` should be left EMPTY. nginx omits an empty header, and a
* WRONG key is a 401 rather than a downgrade to anonymous — verified against the
* live service. A stale token in the environment would therefore break every
* upload while looking like a service fault.
*
* Only the LIST read (`GET /api/uploads/catalog`) still wants a credential;
* reading one batch by its id does not, because the id is itself the proof of
* having sent it.
*/
/**
* Optional-chained because `import.meta.env` is Vite's, and it is undefined
* anywhere Vite is not — the `node --test` runner included. Without the `?.`
* this line throws on import, so every test that so much as names this module
* fails before it runs, with a TypeError that points here rather than at the
* test. Cheap insurance for a value that already has a fallback.
*/
const INGEST_BASE = import.meta.env?.['VITE_INGEST_BASE'] ?? '/ingest';
const ROOT = '/api/uploads/catalog';
/* ── Limits, mirroring the service's own ──────────────────────────────────── */
/**
* 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, from the documented format list. */
export const ACCEPTED_EXTENSIONS = ['.xlsx', '.xls', '.csv', '.tsv'];
/* ── Response types, from the owning team's documented output ─────────────── */
/**
* 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 =
/**
* Accepted and staged, but NOTHING RUNS until an admin releases it.
*
* A review inbox now sits in front of the pipeline — the service answers
* `"Waiting for review. Nothing runs until an admin starts it."` — and this
* status was not in the contract we were given. It matters far more than an
* extra enum member: `isSettled` originally read "not queued and not
* running", so `pending` counted as FINISHED and the panel rendered a
* completed batch reporting nothing imported. An upload that had not yet
* begun would have been shown as a successful import of zero products.
*/
| 'pending'
/**
* Every file in this DROP has been released or dismissed — the drop is spent.
*
* Not an outcome of its own: the answer is on the files. A released file
* carries `released_to`, which is where the run actually is; a dismissed one
* carries nothing because nothing will come. Treating `retired` as finished
* would report a drop that was accepted and is running right now as a
* completed import of zero products.
*/
| 'retired'
| 'queued'
| 'running'
| 'done'
| 'partial'
| 'failed'
| 'interrupted'
| 'cancelled';
/**
* A file inside a drop.
*
* `released` and `dismissed` are the review inbox's two outcomes and neither is
* a result: released means an admin accepted it and the RUN is somewhere else —
* follow `released_to` — while dismissed means they declined it and nothing will
* ever come. Reading either as a finished import reports products that were
* never written.
*/
export type BatchFileStatus =
| 'queued'
| 'running'
| 'done'
| 'failed'
| 'released'
| 'dismissed';
/**
* One product the pipeline wrote, from the run's manifest.
*
* `image_id` is the join key and the only safe one. The owning team calls it
* "the primary key every other product is deduplicated on", and warns that a
* product name differing by one character is a different product — so matching
* a manifest on NAME silently creates duplicates instead of updating.
*
* `unchanged` rows are included on purpose: re-sending a sheet writes nothing,
* and omitting them would make a completely successful upload return an empty
* list that reads as total failure.
*/
export interface IngestProduct {
image_id: string;
brand: string;
product_name: string;
product_sku?: string;
/** `sheet` when the sheet supplied it, `Internal` when the pipeline minted one. */
sku_source?: string;
disposition: 'inserted' | 'backfilled' | 'unchanged';
}
/** 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;
/**
* What the run actually wrote, product by product. Returned on the
* single-batch read only — the list endpoints omit it, because twenty runs of
* thousands of rows is not a list payload.
*/
products?: IngestProduct[];
/** True when the manifest was capped at 5,000 rows for this file. */
products_truncated?: boolean;
}
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;
/**
* The RUN this file became once an admin released it.
*
* Null while it waits and after it is dismissed. The drop id stays valid for
* the whole lifecycle — an earlier build deleted the drop on release and the
* poll started 404ing, which made running, declined and lost look identical
* from outside.
*/
released_to?: 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;
backfilled: number;
skipped_existing: number;
rejected: number;
}
export interface IngestBatch {
batch_id: string;
status: BatchStatus;
detail: string | null;
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 {
readonly status: number;
readonly body: string;
constructor(message: string, status: number, body = '') {
super(message);
this.name = 'IngestError';
this.status = status;
this.body = body;
}
}
/* ── Requests ─────────────────────────────────────────────────────────────── */
export interface SubmitOptions {
files: File[];
/**
* A label for the review inbox, so the admin can see who sent what.
*
* Free text, trimmed to 60 characters by the service, and defaulting to
* "anonymous" when omitted. It is worth sending: the drop endpoint takes no
* credential, so without this every submission in the inbox is indistinguishable
* and an admin approving one cannot tell whose it is.
*/
sender?: string;
signal?: AbortSignal;
}
/**
* 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 is ${files.length} files. The service takes ${MAX_FILES} per upload — send them in smaller batches.`,
413,
);
}
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,
);
}
}
/**
* 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, sender = 'nearle-console', signal } = options;
guardFiles(files);
const form = new FormData();
for (const file of files) form.append('files', file, file.name);
// Labels the drop in the review inbox. The endpoint takes no credential, so
// without this every submission arrives as "anonymous" and the admin deciding
// whether to run it cannot tell ours from anyone else's.
form.append('sender', sender);
// `use_llm` and `fetch_images` are no longer sent, and passing them is inert.
//
// They decide how a run behaves and commit the host to outbound work — image
// search is minutes per batch on one vCPU — so the choice belongs to the admin
// pressing Start, not to whoever dropped the file. Keeping them in the request
// would have read like control we do not have.
let response: Response;
try {
response = await fetch(`${INGEST_BASE}${ROOT}`, {
method: 'POST',
body: form,
// Content-Type is deliberately unset: the browser adds it WITH the
// multipart boundary. Setting it by hand omits the boundary and the
// server parses nothing.
headers: { Accept: 'application/json' },
...(signal ? { signal } : {}),
});
} catch (cause) {
throw new IngestError(
cause instanceof DOMException && cause.name === 'AbortError'
? 'Cancelled.'
: 'Could not reach the ingest service.',
0,
);
}
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 when the batch is sitting in the review inbox, untouched.
*
* Not a failure and not a result — it is waiting for a person. The distinction
* has to be explicit, because the two obvious ways to classify it are both
* wrong: called finished, the screen reports an import of zero products that
* never ran; called in-progress, the browser polls indefinitely for something
* only an admin can move.
*/
export function isAwaitingReview(batch: IngestBatch): boolean {
return batch.status === 'pending' && !releasedRunId(batch);
}
/**
* The run a released drop became, if an admin has accepted it.
*
* A drop is a submission, not a run. Releasing it starts a separate batch and
* records its id on the file as `released_to`; the drop id keeps working and
* keeps saying `released`, so the results are one hop away rather than at the
* id you already hold.
*
* Read off the files rather than the drop, because that is where the service
* puts it — a drop of several files can in principle be released in parts.
*/
export function releasedRunId(batch: IngestBatch): string | null {
for (const file of batch.files ?? []) {
if (file.released_to) return file.released_to;
}
return null;
}
/**
* True when an admin declined the drop. Nothing further will ever arrive, so a
* client that keeps polling is waiting for something that cannot happen.
*/
export function isDismissed(batch: IngestBatch): boolean {
const files = batch.files ?? [];
return files.length > 0 && files.every((file) => file.status === 'dismissed');
}
/**
* Follows a drop to its run, once, and returns whichever is the real answer.
*
* The caller polls a drop id. If it is released, the numbers it wants are on
* the RUN — so this hops and returns that instead. Everything else comes back
* unchanged, so a caller never has to know a drop and a run are different
* things.
*/
export async function resolveBatch(batch: IngestBatch, signal?: AbortSignal): Promise<IngestBatch> {
const runId = releasedRunId(batch);
if (!runId || runId === batch.batch_id) return batch;
try {
return await fetchBatch(runId, signal);
} catch {
// The drop is still the honest answer if the run cannot be read — better a
// stale "released" than an error for something that did succeed.
return batch;
}
}
/**
* The products a finished run wrote, optionally narrowed to our own files.
*
* `filenames` is not optional in practice and should always be passed. An admin
* can assemble ONE run from several drops — the owning team's own words: "a run
* an admin assembled from several drops lists every file in it, so you may see
* filenames batched alongside your own" — so a run's manifest can contain other
* senders' products.
*
* Reading all of them was a real hazard, not a tidiness point. The sheet's price
* and opening stock are applied to whatever the manifest is matched against, so
* a product from someone else's sheet sharing a name with one of our rows would
* have been priced and stocked from OUR file, into OUR merchant's branch.
*
* Filtering by filename is the best this contract allows and it is not airtight:
* two senders can both upload `products.csv`. Narrowing by drop would be exact,
* and the run's files carry no drop reference to narrow by — worth asking for.
*/
export function productsOf(batch: IngestBatch, filenames?: readonly string[]): IngestProduct[] {
const wanted = filenames ? new Set(filenames) : null;
return (batch.files ?? [])
.filter((file) => !wanted || wanted.has(file.filename))
.flatMap((file) => file.result?.products ?? []);
}
/**
* True once the batch has stopped moving, whatever the outcome.
*
* Listed positively rather than as "not queued and not running". The negative
* form silently absorbed every status added later — which is exactly how
* `pending` came to read as a completed import the day the review inbox
* appeared. A new status now shows up as "not settled" and stalls a spinner,
* which is visible, rather than as "done" and fabricates a result.
*/
export function isSettled(batch: IngestBatch): boolean {
// A retired drop whose files went nowhere we can follow is over. Normally
// resolveBatch has already hopped to the run, or isDismissed has caught a
// decline — this is the remainder, and leaving it unsettled would spin a
// progress bar on a drop that no longer exists.
if (batch.status === 'retired') {
return !releasedRunId(batch);
}
return (
batch.status === 'done' ||
batch.status === 'partial' ||
batch.status === 'failed' ||
batch.status === 'interrupted' ||
batch.status === 'cancelled'
);
}
/**
* 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);
// Follows a released drop to the run it became, so the caller polls the
// thing that actually has progress on it rather than a record that will say
// "released" forever.
const batch = await resolveBatch(await fetchBatch(batchId, signal), signal);
onTick(batch);
// Stops on a review hold and on a dismissal as well as on a result. Waiting
// for an admin is not progress, a declined drop will never produce one, and
// a browser tab cannot outlast either — the drop id is what the operator
// comes back with.
if (isSettled(batch) || isAwaitingReview(batch) || isDismissed(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;
try {
payload = text ? JSON.parse(text) : null;
} catch {
payload = text;
}
if (!response.ok) throw describe(response.status, payload, text);
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.
*
* Each of these has one cause and one remedy, and a generic "request failed"
* sends people to look at their spreadsheet for a problem that is in the
* deployment.
*/
function describe(status: number, payload: unknown, text: string): IngestError {
const detail = detailOf(payload);
const body = text.slice(0, 2000);
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 ??
'That key authenticated but does not hold `upload_catalog`. It needs the `uploader` or `admin` role.',
status,
body,
);
}
if (status === 413) {
return new IngestError(
detail ?? 'Too large for the service — 20 files, 10 MB each, 50 MB and 20,000 rows per upload.',
status,
body,
);
}
if (status === 429) {
return new IngestError(
detail ??
'The review inbox is full, so NOTHING was stored — this upload was not merely delayed. An admin has to clear it before you resend.',
status,
body,
);
}
if (status === 400) {
return new IngestError(
detail ??
'The service could not read that file. Every sheet needs a product-name column — product, item, variant or name.',
status,
body,
);
}
return new IngestError(detail ?? `The ingest service returned HTTP ${status}.`, status, body);
}
/* ── Reading a finished batch ─────────────────────────────────────────────── */
/** 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(batch: IngestBatch): string {
const { totals } = batch;
if (isDismissed(batch)) {
// A refusal, not a failure, and nothing further is coming.
//
// The DROP-level detail is deliberately not used here. It still reads
// "Waiting for review. Nothing runs until an admin starts it." on a drop
// that has since been declined — the sentence was written when the file was
// accepted and nothing rewrites it. Rendering it would tell the operator to
// keep waiting for a decision that has already been made against them.
//
// A reason attached to the FILE is the admin's own and is worth showing.
const reason = (batch.files ?? []).map((file) => file.detail).find(Boolean);
return reason
? `An admin declined this upload: ${reason}`
: 'An admin declined this upload. Nothing was imported.';
}
if (isAwaitingReview(batch)) {
// The service's own sentence when it has one — it is clearer than anything
// invented here, and it changes if their review policy does.
return (
batch.detail ??
'Waiting for review. Nothing runs until an admin on the ingest service starts it.'
);
}
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 = [`${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 };
}