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 [recent, setRecent] = useState([]); 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.'); } }; // Runs that started without anyone here pressing anything. // // UPLOAD_AUTORUN sends a colleague's upload straight into the pipeline, so // the inbox it used to wait in is empty and the run has an id this screen was // never told. Without this list an auto-started batch is invisible: its // eleven stages are recorded and served, and nothing renders them. Slim // responses, so twenty of these do not carry twenty product manifests. const loadRecent = async () => { try { const { batches } = await api.listCatalogBatches(20); setRecent(batches || []); } catch { /* the inbox call surfaces auth/network trouble; one banner is enough */ } }; const watch = (run) => { setBatch(run); batchIdRef.current = TERMINAL.has(run.status) ? null : run.batch_id; }; useEffect(() => { loadInbox(); loadRecent(); const timer = setInterval(() => { loadInbox(); loadRecent(); }, 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. Those uploads now start on their own — pick one under Recent runs to follow its stages. Anything still waiting for a decision appears below, where you can tick files 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 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}
)}
{/* ---- runs that started without anyone pressing anything -------- */} {recent.length > 0 && (

Recent runs

Uploads start on arrival, so this is where a colleague's files show up. Pick one to see its eleven stages.

)} {/* ---- the drops waiting to be started --------------------------- */} {inbox && submissions.length === 0 && !batch && (

Nothing waiting for review — files sent to{' '} /api/uploads/catalog start on their own and appear under Recent runs.

Set UPLOAD_AUTORUN=false to hold them here for a decision instead.

)} {submissions.map((submission) => { const ids = submission.files.map((f) => f.file_id); const allOn = ids.length > 0 && ids.every((id) => selected.has(id)); return (

from {submission.submitted_by} {ago(submission.created_at)} · {submission.files.length} file(s)

); })} {/* ---- the run ---------------------------------------------------- */} {batch && (

{batch.status === 'done' && } {batch.status === 'partial' && } {batch.status === 'failed' && } {batch.status === 'running' && } Run — {batch.files_done}/{batch.files_total} file(s) done

{batch.runner} {batch.status}

{batch.batch_id}

{awaitingOrchestrator && (

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.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 (
  • {file.filename} {file.started_at && file.finished_at && ( {fmtSecs(file.started_at, file.finished_at)} )} {file.status}
    {file.detail &&

    {file.detail}

    } {names.length > 0 && (
      {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 (
    1. {index}/11 {done && } {active && } {!record && } {name} {record && record.rows_total > 0 && ( {record.rows_done}/{record.rows_total} )} {done && ( {fmtSecs(record.started_at, record.finished_at)} )}
    2. ); })}
    )}
  • ); } 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 (

    {label}

    {value ?? 0}

    ); }