diff --git a/src/api/client.js b/src/api/client.js
index 2139321..49add64 100644
--- a/src/api/client.js
+++ b/src/api/client.js
@@ -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) =>
diff --git a/src/pages/AdminPage.jsx b/src/pages/AdminPage.jsx
index 712763b..c3e5c1a 100644
--- a/src/pages/AdminPage.jsx
+++ b/src/pages/AdminPage.jsx
@@ -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' && }
+
{activeTab === 'project' && (
diff --git a/src/pages/BatchCatalogPanel.jsx b/src/pages/BatchCatalogPanel.jsx
index a656a03..be7a89f 100644
--- a/src/pages/BatchCatalogPanel.jsx
+++ b/src/pages/BatchCatalogPanel.jsx
@@ -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);
diff --git a/src/pages/InboxPanel.jsx b/src/pages/InboxPanel.jsx
index 1a84863..ccb9665 100644
--- a/src/pages/InboxPanel.jsx
+++ b/src/pages/InboxPanel.jsx
@@ -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);
diff --git a/src/pages/OrchestrationPanel.jsx b/src/pages/OrchestrationPanel.jsx
new file mode 100644
index 0000000..9417600
--- /dev/null
+++ b/src/pages/OrchestrationPanel.jsx
@@ -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 (
+
+
+
+
+
+ Dagster Orchestration
+
+
+ Spreadsheets sent in by a colleague through the upload API. Tick the files you
+ want and press Start selected to send them through the 11-stage
+ pipeline as one batch, run by Dagster instead of this server's own worker.
+ Every stage each file passes through is recorded below and kept after the run
+ finishes.
+
+
+
+
+
+
+
+
+
+ {pending} file(s) awaiting review
+
+
+
+ {/* Same defaults and the same warnings as the other ingestion tabs. */}
+
+
+
+
+
+ {notice && (
+
+ {notice}
+
+ )}
+ {error && (
+
+ {error}
+
+ )}
+
+
+ {/* ---- the drops waiting to be started --------------------------- */}
+ {inbox && submissions.length === 0 && !batch && (
+
+
+
+ Nothing waiting. Files sent to /api/uploads/catalog{' '}
+ appear here to be orchestrated.
+
+
+
+ Staged and waiting for the Dagster orchestrator to claim it. That happens
+ within a minute of batch_upload_sensor{' '}
+ running — if Dagster is not up, which is the case in production, this
+ batch will wait indefinitely.
+
+
+
+
+ )}
+
+
+ {batch.files.map((file) => (
+
+ ))}
+
+
+ {batch.totals && (
+
+
+
+
+
+
+
+ )}
+
+ )}
+
+ {!inbox && (
+
+
+
+ )}
+
+ {running && (
+
+
+ A Dagster run reports through the same manifest this page reads, flushed every few
+ seconds — so progress here can lag the Dagster UI at :3030 by a tick or two.
+
+ )}
+
+ );
+}
+
+/* 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 (
+