Frontend Feature add on Dagster

This commit is contained in:
sriram
2026-08-29 14:51:10 +05:30
parent 8efff41226
commit 80daf30a42
5 changed files with 554 additions and 5 deletions

View File

@@ -385,10 +385,17 @@ export const api = {
// calls - a colleague hits it with an X-API-Key from a script.
listInbox: () => request('/api/admin/catalog-batch/inbox'),
startBatchFromInbox: (fileIds, { use_llm = false, fetch_images = false } = {}) =>
// `runner` picks the executor: 'inprocess' (this container's worker thread,
// the default and the only one that exists in production) or 'dagster'
// (staged and left for the orchestrator to claim). See InboxStartRequest in
// app/api/routers/batch_catalog.py.
startBatchFromInbox: (
fileIds,
{ use_llm = false, fetch_images = false, runner = 'inprocess' } = {},
) =>
request('/api/admin/catalog-batch/from-inbox', {
method: 'POST',
body: JSON.stringify({ file_ids: fileIds, use_llm, fetch_images }),
body: JSON.stringify({ file_ids: fileIds, use_llm, fetch_images, runner }),
}),
dismissInboxFiles: (fileIds) =>

View File

@@ -1,10 +1,11 @@
import React, { useEffect, useState } from 'react';
import { Wrench, Database, UploadCloud, Layers, Inbox } from 'lucide-react';
import { Wrench, Database, UploadCloud, Layers, Inbox, Workflow } from 'lucide-react';
import { api } from '../api/client';
import { NavigationHeader } from '../components/NavigationHeader';
import { StoreCatalogPanel } from './StoreCatalogPanel';
import { BatchCatalogPanel } from './BatchCatalogPanel';
import { InboxPanel } from './InboxPanel';
import { OrchestrationPanel } from './OrchestrationPanel';
/*
* The admin panel carries two views: the read-only project overview, and the
@@ -22,6 +23,7 @@ const TABS = [
{ id: 'ingest', label: 'Store Catalog Ingestion', icon: UploadCloud },
{ id: 'batch', label: 'Batch Catalog Ingestion', icon: Layers },
{ id: 'inbox', label: 'Review Inbox', icon: Inbox },
{ id: 'orchestration', label: 'Dagster Orchestration', icon: Workflow },
];
// How often the pending count is refreshed. Slower than the 3s a running batch
@@ -127,6 +129,11 @@ export function AdminPage() {
/>
)}
{/* Starts from the same inbox but hands the batch to Dagster, and stays
put to show the run rather than handing off to the Batch tab - the
stage-by-stage view is the reason to be on this tab at all. */}
{activeTab === 'orchestration' && <OrchestrationPanel />}
{activeTab === 'project' && (
<div className="space-y-6">
<div className="rounded-2xl border border-ink-900/10 bg-paper-50 p-6 shadow-xs">

View File

@@ -58,7 +58,11 @@ export function BatchCatalogPanel({ adoptBatch = null }) {
const [error, setError] = useState('');
const [batch, setBatch] = useState(null);
const [useLlm, setUseLlm] = useState(false);
const [fetchImages, setFetchImages] = useState(false);
// Ticked by default: a batch ingested without it produces rows with no
// images at all, and nothing downstream ever goes back to fill them in. The
// API default stays off - an anonymous upload must not trigger network
// calls - so this is the admin, in front of the checkbox, opting in.
const [fetchImages, setFetchImages] = useState(true);
const pollRef = useRef(null);
const batchIdRef = useRef(null);

View File

@@ -46,7 +46,10 @@ export function InboxPanel({ onBatchStarted }) {
const [error, setError] = useState('');
const [notice, setNotice] = useState('');
const [useLlm, setUseLlm] = useState(false);
const [fetchImages, setFetchImages] = useState(false);
// Ticked by default - see the same note in BatchCatalogPanel. A file
// released from the inbox is the main way catalog rows are created, and
// without this they arrive with no images.
const [fetchImages, setFetchImages] = useState(true);
const pollRef = useRef(null);

View File

@@ -0,0 +1,528 @@
import React, { useEffect, useRef, useState } from 'react';
import {
AlertTriangle, ArrowRight, CheckCircle2, Clock, FileSpreadsheet,
Loader2, MinusCircle, PlayCircle, RefreshCw, Workflow, XCircle,
} from 'lucide-react';
import { api } from '../api/client';
/*
* Admin -> Dagster Orchestration.
*
* Same starting point as the Review Inbox: spreadsheets a colleague sent
* through the upload API, ticked and started as one batch. Two things are
* different, and they are the reason this is its own screen.
*
* It hands the batch to DAGSTER rather than to this container's worker thread.
* Both run the identical eleven stages out of app/core/batch_ingest - Dagster
* adds lineage, retries and a re-run launchpad on a development machine. Which
* one owns a batch is recorded on the batch itself (`runner`), because both
* watch the same directory and, before that field existed, enabling the sensor
* beside a live API meant both ingested every batch twice.
*
* And it shows the whole pipeline, not just the current position. The Batch tab
* prints one line - "Stage 5/11 ..." - which disappears the moment a file
* finishes, so you can never see what a completed file actually did. Here every
* file carries its timeline: which stages ran, how long each took, how many rows
* each saw, kept after the run ends.
*
* Progress is read from the ordinary batch endpoints, NOT from Dagster. A run
* flushes its manifest to disk every few seconds and the API serves batches from
* that same manifest, so this screen follows a Dagster run in another process
* with no coupling between the two - and keeps working in production, where
* Dagster is not deployed at all.
*/
// The inbox half asks "has anything arrived?"; the run half asks "how far has
// it got?". Different questions, different cadences - and the run poll stops
// itself at a terminal status so a tab left open overnight goes quiet.
const INBOX_POLL_MS = 5000;
const RUN_POLL_MS = 3000;
const TERMINAL = new Set(['done', 'failed', 'partial', 'cancelled']);
const FILE_TONE = {
done: 'bg-emerald-500/10 text-emerald-700 border-emerald-500/20',
failed: 'bg-maroon-100 text-maroon-600 border-maroon-500/30',
running: 'bg-amber-500/10 text-amber-700 border-amber-500/20',
queued: 'bg-white text-slate-500 border-ink-900/10',
cancelled: 'bg-ink-100 text-slate-500 border-ink-900/10',
};
const FILE_ICON = {
done: CheckCircle2,
failed: XCircle,
running: Loader2,
queued: Clock,
cancelled: MinusCircle,
};
function ago(seconds) {
const s = Math.max(0, Math.floor(Date.now() / 1000 - seconds));
if (s < 60) return 'just now';
if (s < 3600) return `${Math.floor(s / 60)}m ago`;
if (s < 86400) return `${Math.floor(s / 3600)}h ago`;
return `${Math.floor(s / 86400)}d ago`;
}
function fmtBytes(n) {
if (!n) return '';
if (n < 1024) return `${n} B`;
if (n < 1024 * 1024) return `${(n / 1024).toFixed(0)} KB`;
return `${(n / (1024 * 1024)).toFixed(1)} MB`;
}
function fmtSecs(from, to) {
if (!from || !to) return '';
const s = to - from;
if (s < 1) return '<1s';
if (s < 60) return `${s.toFixed(1)}s`;
return `${Math.floor(s / 60)}m ${Math.round(s % 60)}s`;
}
export function OrchestrationPanel() {
const [inbox, setInbox] = useState(null);
const [selected, setSelected] = useState(() => new Set());
const [batch, setBatch] = useState(null);
const [busy, setBusy] = useState('');
const [error, setError] = useState('');
const [notice, setNotice] = useState('');
const [useLlm, setUseLlm] = useState(false);
const [fetchImages, setFetchImages] = useState(true);
const batchIdRef = useRef(null);
const loadInbox = async () => {
try {
const next = await api.listInbox();
setInbox(next);
// Drop a selection the server no longer offers - another tab may have
// started or dismissed it, and a checkbox pointing at a file that is
// already gone would 409 on the next click with no explanation.
const live = new Set(next.submissions.flatMap((s) => s.files.map((f) => f.file_id)));
setSelected((current) => new Set([...current].filter((id) => live.has(id))));
} catch (err) {
setError(err?.message || 'Could not read the inbox.');
}
};
useEffect(() => {
loadInbox();
const timer = setInterval(loadInbox, INBOX_POLL_MS);
return () => clearInterval(timer);
// eslint-disable-next-line react-hooks/exhaustive-deps
}, []);
useEffect(() => {
const timer = setInterval(async () => {
const id = batchIdRef.current;
if (!id) return;
try {
const next = await api.getCatalogBatch(id);
setBatch(next);
// A dagster batch sits `queued` until the orchestrator claims it, which
// is not a terminal state - keep watching it, that wait is the thing
// the operator is here to see.
if (TERMINAL.has(next.status)) batchIdRef.current = null;
} catch {
/* a dropped poll is not worth surfacing; the next tick retries */
}
}, RUN_POLL_MS);
return () => clearInterval(timer);
}, []);
const toggle = (fileId) => {
setSelected((current) => {
const next = new Set(current);
if (next.has(fileId)) next.delete(fileId);
else next.add(fileId);
return next;
});
};
const toggleDrop = (submission) => {
const ids = submission.files.map((f) => f.file_id);
const allOn = ids.every((id) => selected.has(id));
setSelected((current) => {
const next = new Set(current);
ids.forEach((id) => (allOn ? next.delete(id) : next.add(id)));
return next;
});
};
const handleStart = async () => {
if (!selected.size) return;
setBusy('start');
setError('');
setNotice('');
try {
const started = await api.startBatchFromInbox([...selected], {
use_llm: useLlm,
fetch_images: fetchImages,
runner: 'dagster',
});
setSelected(new Set());
setBatch(started);
batchIdRef.current = started.batch_id;
setNotice(
`Staged ${started.files_total} file(s) as one batch for Dagster. It will start ` +
`when the orchestrator picks it up.`
);
await loadInbox();
} catch (err) {
setError(err?.message || 'Could not stage those files.');
await loadInbox();
} finally {
setBusy('');
}
};
// The way out of a batch nothing came for. Resume hands it to this
// container's worker AND takes ownership, so Dagster will not also claim it.
const handleRunHere = async () => {
if (!batch) return;
setBusy('resume');
setError('');
try {
const resumed = await api.resumeCatalogBatch(batch.batch_id);
setBatch(resumed);
batchIdRef.current = resumed.batch_id;
setNotice('Running this batch here instead. Dagster will no longer pick it up.');
} catch (err) {
setError(err?.message || 'Could not run the batch here.');
} finally {
setBusy('');
}
};
const submissions = inbox?.submissions || [];
const pending = inbox?.pending_count ?? 0;
const running = batch && !TERMINAL.has(batch.status);
// Queued AND owned by Dagster means nothing has claimed it: either the
// orchestrator is not running, or its sensor has not ticked yet.
const awaitingOrchestrator =
batch && batch.runner === 'dagster' && batch.status === 'queued';
return (
<div className="space-y-6">
<div className="rounded-2xl border border-ink-900/10 bg-paper-50 p-6 shadow-xs">
<div className="flex items-start justify-between gap-3">
<div>
<h2 className="font-display text-lg font-bold text-ink-950 flex items-center gap-2 mb-2">
<Workflow className="h-5 w-5 text-amber-500" /> Dagster Orchestration
</h2>
<p className="text-xs text-slate-500 leading-relaxed max-w-3xl">
Spreadsheets sent in by a colleague through the upload API. Tick the files you
want and press <strong>Start selected</strong> to send them through the 11-stage
pipeline as one batch, run by Dagster instead of this server&apos;s own worker.
Every stage each file passes through is recorded below and kept after the run
finishes.
</p>
</div>
<button
type="button"
onClick={loadInbox}
aria-label="Refresh the inbox"
className="shrink-0 flex items-center gap-1.5 px-3 py-1.5 rounded-lg border border-ink-900/15 bg-white text-[11px] font-bold text-slate-600 hover:text-ink-950 transition cursor-pointer"
>
<RefreshCw className="h-3.5 w-3.5" /> Refresh
</button>
</div>
<div className="mt-5 flex flex-wrap items-center gap-3">
<button
type="button"
onClick={handleStart}
disabled={!selected.size || busy}
className="flex items-center gap-2 px-4 py-2 rounded-xl bg-amber-500 text-slate-950 text-xs font-bold hover:bg-amber-600 disabled:opacity-40 transition cursor-pointer"
>
{busy === 'start'
? <Loader2 className="h-4 w-4 animate-spin" />
: <PlayCircle className="h-4 w-4" />}
Start selected{selected.size ? ` (${selected.size})` : ''}
<ArrowRight className="h-3.5 w-3.5" />
</button>
<span className="text-[11px] text-slate-400 font-mono">
{pending} file(s) awaiting review
</span>
</div>
{/* Same defaults and the same warnings as the other ingestion tabs. */}
<div className="mt-4 flex flex-col gap-2">
<label className="flex items-start gap-2 text-[11px] text-slate-600 cursor-pointer">
<input
type="checkbox"
checked={fetchImages}
onChange={(e) => setFetchImages(e.target.checked)}
className="mt-0.5 accent-amber-500 cursor-pointer"
/>
<span>
<span className="font-bold text-ink-950">Search for product images</span>
{' '}&mdash; much slower, and makes outbound requests for every row.
</span>
</label>
<label className="flex items-start gap-2 text-[11px] text-slate-600 cursor-pointer">
<input
type="checkbox"
checked={useLlm}
onChange={(e) => setUseLlm(e.target.checked)}
className="mt-0.5 accent-amber-500 cursor-pointer"
/>
<span>
<span className="font-bold text-ink-950">Use the LLM for missing descriptions</span>
{' '}&mdash; no effect unless Ollama is reachable, which it is not in production.
</span>
</label>
</div>
{notice && (
<div className="mt-4 rounded-lg bg-emerald-500/5 border border-emerald-500/20 p-3 text-xs text-emerald-700 flex items-start gap-2">
<CheckCircle2 className="h-4 w-4 shrink-0 mt-0.5" /> {notice}
</div>
)}
{error && (
<div className="mt-4 rounded-lg bg-maroon-100 border border-maroon-500/30 p-3 text-xs text-maroon-600 flex items-start gap-2">
<XCircle className="h-4 w-4 shrink-0 mt-0.5" /> {error}
</div>
)}
</div>
{/* ---- the drops waiting to be started --------------------------- */}
{inbox && submissions.length === 0 && !batch && (
<div className="rounded-2xl border border-ink-900/10 bg-paper-50 p-10 text-center">
<Workflow className="h-8 w-8 text-slate-300 mx-auto mb-3" />
<p className="text-xs text-slate-500">
Nothing waiting. Files sent to <code className="font-mono">/api/uploads/catalog</code>{' '}
appear here to be orchestrated.
</p>
</div>
)}
{submissions.map((submission) => {
const ids = submission.files.map((f) => f.file_id);
const allOn = ids.length > 0 && ids.every((id) => selected.has(id));
return (
<div
key={submission.submission_id}
className="rounded-2xl border border-ink-900/10 bg-paper-50 p-6 shadow-xs"
>
<div className="flex items-center justify-between gap-3 mb-4">
<h3 className="font-display text-sm font-bold text-ink-950">
from <span className="text-amber-600">{submission.submitted_by}</span>
<span className="ml-2 font-normal text-slate-400">
{ago(submission.created_at)} &middot; {submission.files.length} file(s)
</span>
</h3>
<button
type="button"
onClick={() => toggleDrop(submission)}
className="shrink-0 text-[11px] font-bold text-slate-500 hover:text-ink-950 transition cursor-pointer"
>
{allOn ? 'Clear all' : 'Select all'}
</button>
</div>
<ul className="space-y-1.5">
{submission.files.map((file) => (
<li key={file.file_id}>
<label className="flex items-center gap-3 px-3 py-2.5 rounded-xl border border-ink-900/10 bg-white cursor-pointer hover:border-amber-500/40 transition">
<input
type="checkbox"
checked={selected.has(file.file_id)}
onChange={() => toggle(file.file_id)}
className="accent-amber-500 cursor-pointer"
/>
<FileSpreadsheet className="h-4 w-4 text-amber-600 shrink-0" />
<span className="text-[11px] font-bold text-ink-950 truncate flex-1">
{file.filename}
</span>
<span className="font-mono text-[10px] text-slate-400 shrink-0">
{file.rows_total} rows
</span>
<span className="font-mono text-[10px] text-slate-300 shrink-0 w-16 text-right">
{fmtBytes(file.size_bytes)}
</span>
</label>
</li>
))}
</ul>
</div>
);
})}
{/* ---- the run ---------------------------------------------------- */}
{batch && (
<div className="rounded-2xl border border-ink-900/10 bg-paper-50 p-6 shadow-xs">
<div className="flex items-center justify-between gap-3 mb-1">
<h3 className="font-display text-sm font-bold text-ink-950 flex items-center gap-2">
{batch.status === 'done' && <CheckCircle2 className="h-4 w-4 text-leaf-600" />}
{batch.status === 'partial' && <AlertTriangle className="h-4 w-4 text-amber-600" />}
{batch.status === 'failed' && <XCircle className="h-4 w-4 text-maroon-600" />}
{batch.status === 'running' && <Loader2 className="h-4 w-4 animate-spin text-amber-600" />}
Run &mdash; {batch.files_done}/{batch.files_total} file(s) done
</h3>
<div className="flex items-center gap-2 shrink-0">
<span className="inline-flex items-center gap-1 rounded-full bg-amber-500/15 border border-amber-500/30 px-2.5 py-0.5 text-[10px] font-bold text-amber-600">
<Workflow className="h-3 w-3" /> {batch.runner}
</span>
<span className="font-mono text-[11px] uppercase text-slate-400">{batch.status}</span>
</div>
</div>
<p className="font-mono text-[10px] text-slate-400 mb-4">{batch.batch_id}</p>
{awaitingOrchestrator && (
<div className="mb-4 rounded-lg bg-amber-500/5 border border-amber-500/20 p-3 text-[11px] text-amber-700">
<p className="flex items-start gap-2">
<Clock className="h-4 w-4 shrink-0 mt-0.5" />
<span>
Staged and waiting for the Dagster orchestrator to claim it. That happens
within a minute of <code className="font-mono">batch_upload_sensor</code>{' '}
running &mdash; if Dagster is not up, which is the case in production, this
batch will wait indefinitely.
</span>
</p>
<button
type="button"
onClick={handleRunHere}
disabled={busy === 'resume'}
className="mt-2.5 flex items-center gap-1.5 px-3 py-1.5 rounded-lg bg-amber-500 text-slate-950 text-[11px] font-bold hover:bg-amber-600 disabled:opacity-40 transition cursor-pointer"
>
{busy === 'resume'
? <Loader2 className="h-3.5 w-3.5 animate-spin" />
: <PlayCircle className="h-3.5 w-3.5" />}
Run here instead
</button>
</div>
)}
<ul className="space-y-3">
{batch.files.map((file) => (
<FileStages
key={`${file.index}-${file.filename}`}
file={file}
stageNames={batch.stage_names || []}
/>
))}
</ul>
{batch.totals && (
<div className="mt-5 grid grid-cols-2 sm:grid-cols-3 lg:grid-cols-5 gap-3">
<Stat label="Rows read" value={batch.totals.rows_total} tone="slate" />
<Stat label="Products built" value={batch.totals.products_built} tone="blue" />
<Stat label="Inserted" value={batch.totals.inserted} tone="emerald" />
<Stat label="Backfilled" value={batch.totals.backfilled} tone="amber" />
<Stat label="Row errors" value={batch.totals.error_count} tone="maroon" />
</div>
)}
</div>
)}
{!inbox && (
<div className="rounded-2xl border border-ink-900/10 bg-paper-50 p-10 text-center">
<Loader2 className="h-5 w-5 animate-spin text-amber-500 mx-auto" />
</div>
)}
{running && (
<p className="text-[11px] text-slate-400 flex items-start gap-1.5">
<AlertTriangle className="h-3.5 w-3.5 shrink-0 mt-0.5" />
A Dagster run reports through the same manifest this page reads, flushed every few
seconds &mdash; so progress here can lag the Dagster UI at :3030 by a tick or two.
</p>
)}
</div>
);
}
/* One file, and every stage of the pipeline as it happened to that file.
*
* All eleven are drawn from the moment the batch exists, greyed out until
* reached, so the shape of the pipeline is visible before anything runs and
* does not reflow as stages appear.
*
* A stage is looked up BY INDEX rather than by position in `file.stages`.
* Stages 8-11 run once per brand in the sheet, so the backend folds repeat
* visits into one record per index; matching on index is what keeps a
* three-brand file from appearing to run backwards. */
function FileStages({ file, stageNames }) {
const Icon = FILE_ICON[file.status] || Clock;
const byIndex = new Map((file.stages || []).map((s) => [s.index, s]));
const names = stageNames.length ? stageNames : (file.stages || []).map((s) => s.name);
return (
<li className={`rounded-xl border p-3 ${FILE_TONE[file.status] || FILE_TONE.queued}`}>
<div className="flex items-center justify-between gap-2">
<span className="flex items-center gap-2 min-w-0">
<Icon className={`h-3.5 w-3.5 shrink-0 ${file.status === 'running' ? 'animate-spin' : ''}`} />
<span className="text-[11px] font-bold truncate">{file.filename}</span>
</span>
<span className="flex items-center gap-2 shrink-0">
{file.started_at && file.finished_at && (
<span className="font-mono text-[10px] opacity-70">
{fmtSecs(file.started_at, file.finished_at)}
</span>
)}
<span className="font-mono text-[10px] uppercase">{file.status}</span>
</span>
</div>
{file.detail && <p className="mt-1.5 text-[10px] leading-relaxed">{file.detail}</p>}
{names.length > 0 && (
<ol className="mt-2.5 space-y-1">
{names.map((name, i) => {
const index = i + 1;
const record = byIndex.get(index);
const done = record && record.finished_at;
const active = record && !record.finished_at;
return (
<li
key={index}
className={`flex items-center gap-2 rounded-lg px-2 py-1 text-[10px] ${
done
? 'bg-white/60 text-slate-600'
: active
? 'bg-amber-500/15 text-amber-700 font-bold'
: 'text-slate-400'
}`}
>
<span className="font-mono w-7 shrink-0 text-right opacity-60">{index}/11</span>
{done && <CheckCircle2 className="h-3 w-3 shrink-0 text-leaf-600" />}
{active && <Loader2 className="h-3 w-3 shrink-0 animate-spin" />}
{!record && <Clock className="h-3 w-3 shrink-0 opacity-40" />}
<span className="truncate flex-1">{name}</span>
{record && record.rows_total > 0 && (
<span className="font-mono shrink-0 opacity-70">
{record.rows_done}/{record.rows_total}
</span>
)}
{done && (
<span className="font-mono shrink-0 w-12 text-right opacity-50">
{fmtSecs(record.started_at, record.finished_at)}
</span>
)}
</li>
);
})}
</ol>
)}
</li>
);
}
const TONES = {
emerald: 'border-emerald-500/20 bg-emerald-500/5 text-emerald-700',
blue: 'border-blue-500/20 bg-blue-500/5 text-blue-700',
slate: 'border-ink-900/10 bg-white text-slate-600',
maroon: 'border-maroon-500/20 bg-maroon-100/40 text-maroon-600',
amber: 'border-amber-500/20 bg-amber-500/5 text-amber-700',
};
function Stat({ label, value, tone }) {
return (
<div className={`rounded-xl border p-3 ${TONES[tone] || TONES.slate}`}>
<p className="text-[10px] font-semibold uppercase">{label}</p>
<p className="font-mono text-2xl font-extrabold text-ink-950 mt-0.5">{value ?? 0}</p>
</div>
);
}