diff --git a/studio/backend/core/data_recipe/jobs/manager.py b/studio/backend/core/data_recipe/jobs/manager.py
index fd8a290071..e5e05b0868 100644
--- a/studio/backend/core/data_recipe/jobs/manager.py
+++ b/studio/backend/core/data_recipe/jobs/manager.py
@@ -147,6 +147,15 @@ class JobManager:
"ok": job.progress.ok,
"failed": job.progress.failed,
},
+ "column_progress": {
+ "done": job.column_progress.done,
+ "total": job.column_progress.total,
+ "percent": job.column_progress.percent,
+ "eta_sec": job.column_progress.eta_sec,
+ "rate": job.column_progress.rate,
+ "ok": job.column_progress.ok,
+ "failed": job.column_progress.failed,
+ },
"model_usage": {
name: {
"model": usage.model,
@@ -207,31 +216,31 @@ class JobManager:
return None
in_memory_dataset = self._job.dataset
artifact_path = self._job.artifact_path
+ job_status = self._job.status
if in_memory_dataset is not None:
total = len(in_memory_dataset)
rows = in_memory_dataset[offset:offset + limit]
return {"dataset": rows, "total": total}
if not artifact_path:
- return None
-
- try:
- from data_designer.engine.dataset_builders.artifact_storage import ArtifactStorage
- except Exception:
+ if job_status in {"completed", "error", "cancelled"}:
+ return {"error": "artifact path missing"}
return None
try:
base_dataset_path = Path(artifact_path)
- storage = ArtifactStorage(
- artifact_path=str(base_dataset_path.parent),
- dataset_name=base_dataset_path.name,
- )
- dataframe = storage.load_dataset()
+ parquet_dir = base_dataset_path / "parquet-files"
+ if not parquet_dir.exists():
+ return {"error": f"dataset path missing: {parquet_dir}"}
+
+ from data_designer.config.utils.io_helpers import read_parquet_dataset
+
+ dataframe = read_parquet_dataset(parquet_dir)
total = int(len(dataframe.index))
rows = dataframe.iloc[offset:offset + limit].to_dict(orient="records")
return {"dataset": _to_jsonable(rows), "total": total}
- except Exception:
- return None
+ except Exception as exc:
+ return {"error": f"dataset load failed: {exc}"}
def subscribe(self, job_id: str, *, after_seq: int | None = None) -> Subscription | None:
"""SSE subscribe: get replay buffer + live events stream."""
@@ -347,6 +356,9 @@ class JobManager:
self._job.artifact_path = event.get("artifact_path")
self._job.dataset = event.get("dataset")
self._job.processor_artifacts = event.get("processor_artifacts")
+ if self._job.progress.total and self._job.progress.total > 0:
+ self._job.progress.done = self._job.progress.total
+ self._job.progress.percent = 100.0
if et == "job.error":
self._job.status = "error"
self._job.finished_at = time.time()
diff --git a/studio/backend/core/data_recipe/jobs/parse.py b/studio/backend/core/data_recipe/jobs/parse.py
index 00c05833f7..4a3ca99000 100644
--- a/studio/backend/core/data_recipe/jobs/parse.py
+++ b/studio/backend/core/data_recipe/jobs/parse.py
@@ -132,12 +132,15 @@ def apply_update(job: Job, update: ParsedUpdate) -> None:
job.stage = update.stage
if update.current_column is not None:
job.current_column = update.current_column
+ if update.stage == "generating" and update.current_column not in job._seen_generation_columns:
+ job._seen_generation_columns.append(update.current_column)
if update.rows is not None:
job.rows = update.rows
if update.cols is not None:
job.cols = update.cols
if update.progress is not None:
- job.progress = update.progress
+ job.column_progress = update.progress
+ job.progress = _compute_overall_progress(job, update.progress)
if update.batch_idx is not None:
job.batch.idx = update.batch_idx
if update.batch_total is not None:
@@ -194,6 +197,48 @@ def apply_update(job: Job, update: ParsedUpdate) -> None:
usage.rpm = update.usage_rpm
+def _compute_overall_progress(job: Job, column_progress: Progress) -> Progress:
+ if not job.rows or not job.current_column:
+ return column_progress
+
+ total_rows = max(1, int(job.rows))
+ total_cols = max(
+ 1,
+ len(job._seen_generation_columns),
+ int(job.cols or 0),
+ )
+ current_done = 0 if column_progress.done is None else int(column_progress.done)
+ current_done = max(0, min(current_done, total_rows))
+
+ try:
+ col_index = job._seen_generation_columns.index(job.current_column)
+ except ValueError:
+ col_index = max(0, len(job._seen_generation_columns) - 1)
+
+ col_index = max(0, min(col_index, total_cols - 1))
+ total = total_rows * total_cols
+ done = min(total, (col_index * total_rows) + current_done)
+ prev_done = int(job.progress.done or 0)
+ if done < prev_done:
+ done = prev_done
+ if done > total:
+ done = total
+ percent = (done / total) * 100 if total > 0 else 100.0
+ prev_percent = float(job.progress.percent or 0.0)
+ if percent < prev_percent:
+ percent = prev_percent
+
+ return Progress(
+ done=done,
+ total=total,
+ percent=percent,
+ eta_sec=column_progress.eta_sec,
+ rate=column_progress.rate,
+ ok=column_progress.ok,
+ failed=column_progress.failed,
+ )
+
+
def coerce_event(obj: Any) -> dict:
# worker sends dict already
return obj if isinstance(obj, dict) else {"type": "log", "message": str(obj)}
diff --git a/studio/backend/core/data_recipe/jobs/types.py b/studio/backend/core/data_recipe/jobs/types.py
index d475b928e1..f532741370 100644
--- a/studio/backend/core/data_recipe/jobs/types.py
+++ b/studio/backend/core/data_recipe/jobs/types.py
@@ -52,6 +52,7 @@ class Job:
stage: str | None = None
current_column: str | None = None
progress: Progress = field(default_factory=Progress)
+ column_progress: Progress = field(default_factory=Progress)
batch: BatchProgress = field(default_factory=BatchProgress)
rows: int | None = None
cols: int | None = None
@@ -66,3 +67,4 @@ class Job:
model_usage: dict[str, ModelUsage] = field(default_factory=dict)
_current_usage_model: str | None = None
_in_usage_summary: bool = False
+ _seen_generation_columns: list[str] = field(default_factory=list)
diff --git a/studio/backend/core/data_recipe/jobs/worker.py b/studio/backend/core/data_recipe/jobs/worker.py
index d3c0da79a0..3f105132ab 100644
--- a/studio/backend/core/data_recipe/jobs/worker.py
+++ b/studio/backend/core/data_recipe/jobs/worker.py
@@ -81,10 +81,14 @@ def run_job_process(
rows = int(run.get("rows") or 1000)
dataset_name = str(run.get("dataset_name") or "dataset")
+ artifact_path_raw = run.get("artifact_path")
+ artifact_path = None
+ if isinstance(artifact_path_raw, str) and artifact_path_raw.strip():
+ artifact_path = artifact_path_raw.strip()
run_config_raw = run.get("run_config") or {}
builder = build_config_builder(recipe)
- designer = create_data_designer(recipe)
+ designer = create_data_designer(recipe, artifact_path=artifact_path)
if run_config_raw:
designer.set_run_config(RunConfig.model_validate(run_config_raw))
diff --git a/studio/backend/core/data_recipe/service.py b/studio/backend/core/data_recipe/service.py
index 8a069cf21b..dca3e867ce 100644
--- a/studio/backend/core/data_recipe/service.py
+++ b/studio/backend/core/data_recipe/service.py
@@ -110,10 +110,15 @@ def build_config_builder(recipe: dict[str, Any]):
return DataDesignerConfigBuilder.from_config({"data_designer": recipe_core})
-def create_data_designer(recipe: dict[str, Any]):
+def create_data_designer(
+ recipe: dict[str, Any],
+ *,
+ artifact_path: str | None = None,
+):
from data_designer.interface.data_designer import DataDesigner
return DataDesigner(
+ artifact_path=artifact_path,
model_providers=build_model_providers(recipe),
mcp_providers=build_mcp_providers(recipe),
)
diff --git a/studio/backend/routes/data_recipe.py b/studio/backend/routes/data_recipe.py
index 48cc7339ee..4f4cfc4290 100644
--- a/studio/backend/routes/data_recipe.py
+++ b/studio/backend/routes/data_recipe.py
@@ -17,8 +17,8 @@ if str(backend_path) not in sys.path:
sys.path.insert(0, str(backend_path))
from core.data_recipe.jobs import get_job_manager
-from core.data_recipe.service import preview_recipe, validate_recipe
-from models.data_recipe import JobCreateResponse, PreviewResponse, RecipePayload, ValidateError, ValidateResponse
+from core.data_recipe.service import validate_recipe
+from models.data_recipe import JobCreateResponse, RecipePayload, ValidateError, ValidateResponse
router = APIRouter()
@@ -47,25 +47,6 @@ def validate(payload: RecipePayload) -> ValidateResponse:
return ValidateResponse(valid=True)
-@router.post("/preview", response_model=PreviewResponse)
-def preview(payload: RecipePayload) -> PreviewResponse:
- recipe = payload.recipe
- if not recipe.get("columns"):
- raise HTTPException(status_code=400, detail="Recipe must include columns.")
-
- run = payload.run or {}
- num_records = int(run.get("rows") or 5)
-
- try:
- dataset, artifacts, analysis = preview_recipe(recipe, num_records)
- except RuntimeError as exc:
- raise HTTPException(status_code=503, detail=str(exc)) from exc
- except Exception as exc:
- raise HTTPException(status_code=400, detail=str(exc)) from exc
-
- return PreviewResponse(dataset=dataset, processor_artifacts=artifacts, analysis=analysis)
-
-
@router.post("/jobs", response_class=JSONResponse, response_model=JobCreateResponse)
def create_job(payload: RecipePayload):
recipe = payload.recipe
@@ -73,6 +54,10 @@ def create_job(payload: RecipePayload):
raise HTTPException(status_code=400, detail="Recipe must include columns.")
run: dict[str, Any] = payload.run or {}
+ execution_type = str(run.get("execution_type") or "full").strip().lower()
+ if execution_type not in {"preview", "full"}:
+ raise HTTPException(status_code=400, detail="invalid execution_type: must be 'preview' or 'full'")
+ run["execution_type"] = execution_type
run_config_raw = run.get("run_config")
if run_config_raw is not None:
try:
@@ -139,6 +124,8 @@ def job_dataset(
result = mgr.get_dataset(job_id, limit=limit, offset=offset)
if result is None:
raise HTTPException(status_code=404, detail="dataset not ready")
+ if "error" in result:
+ raise HTTPException(status_code=422, detail=result["error"])
return {
"dataset": result["dataset"],
"total": result["total"],
diff --git a/studio/frontend/src/features/recipe-studio/api/index.ts b/studio/frontend/src/features/recipe-studio/api/index.ts
index 4bf879e207..f166fc95cd 100644
--- a/studio/frontend/src/features/recipe-studio/api/index.ts
+++ b/studio/frontend/src/features/recipe-studio/api/index.ts
@@ -3,13 +3,6 @@ const DEFAULT_BASE = "/api/data-recipe";
export const DATA_DESIGNER_API_BASE =
import.meta.env.VITE_DATA_DESIGNER_API ?? DEFAULT_BASE;
-export type PreviewResponse = {
- dataset?: unknown[];
- // biome-ignore lint/style/useNamingConvention: api schema
- processor_artifacts?: Record;
- analysis?: Record;
-};
-
export type JobCreateResponse = {
// biome-ignore lint/style/useNamingConvention: api schema
job_id: string;
@@ -37,6 +30,17 @@ export type JobStatusResponse = {
failed?: number | null;
};
// biome-ignore lint/style/useNamingConvention: api schema
+ column_progress?: {
+ done?: number | null;
+ total?: number | null;
+ percent?: number | null;
+ // biome-ignore lint/style/useNamingConvention: api schema
+ eta_sec?: number | null;
+ rate?: number | null;
+ ok?: number | null;
+ failed?: number | null;
+ };
+ // biome-ignore lint/style/useNamingConvention: api schema
model_usage?: Record;
rows?: number | null;
cols?: number | null;
@@ -161,10 +165,6 @@ function parseJobEvent(rawEvent: string): JobEvent | null {
};
}
-export async function previewRecipe(payload: unknown): Promise {
- return postJson("/preview", payload);
-}
-
export async function validateRecipe(
payload: unknown,
): Promise {
diff --git a/studio/frontend/src/features/recipe-studio/components/executions/executions-view.tsx b/studio/frontend/src/features/recipe-studio/components/executions/executions-view.tsx
index 6e737c5c9b..8e729c3402 100644
--- a/studio/frontend/src/features/recipe-studio/components/executions/executions-view.tsx
+++ b/studio/frontend/src/features/recipe-studio/components/executions/executions-view.tsx
@@ -388,6 +388,14 @@ export function ExecutionsView({
ETA: {selectedExecution.progress?.eta_sec ?? "--"} s
+ {selectedExecution.current_column && selectedExecution.column_progress && (
+
+ Column {selectedExecution.current_column}:{" "}
+ {selectedExecution.column_progress.done ?? "--"}/
+ {selectedExecution.column_progress.total ?? "--"} (
+ {formatPercent(selectedExecution.column_progress.percent)})
+
+ )}
)}
diff --git a/studio/frontend/src/features/recipe-studio/execution-types.ts b/studio/frontend/src/features/recipe-studio/execution-types.ts
index 06cc4511e9..3a4d39eb18 100644
--- a/studio/frontend/src/features/recipe-studio/execution-types.ts
+++ b/studio/frontend/src/features/recipe-studio/execution-types.ts
@@ -49,6 +49,8 @@ export type RecipeExecutionRecord = {
current_column: string | null;
progress: RecipeExecutionProgress | null;
// biome-ignore lint/style/useNamingConvention: backend schema
+ column_progress: RecipeExecutionProgress | null;
+ // biome-ignore lint/style/useNamingConvention: backend schema
model_usage: Record | null;
// biome-ignore lint/style/useNamingConvention: backend schema
lastEventId: number | null;
diff --git a/studio/frontend/src/features/recipe-studio/executions/execution-helpers.ts b/studio/frontend/src/features/recipe-studio/executions/execution-helpers.ts
new file mode 100644
index 0000000000..ab8b864006
--- /dev/null
+++ b/studio/frontend/src/features/recipe-studio/executions/execution-helpers.ts
@@ -0,0 +1,163 @@
+import type {
+ RecipeExecutionAnalysis,
+ RecipeExecutionRecord,
+ RecipeExecutionStatus,
+} from "../execution-types";
+import type { RecipePayload } from "../utils/payload/types";
+
+export const DATASET_PAGE_SIZE = 20;
+
+export function buildSignature(name: string, payload: RecipePayload): string {
+ return JSON.stringify({ name, payload });
+}
+
+export function formatSavedLabel(savedAt: number | null): string {
+ if (!savedAt) {
+ return "Not saved yet";
+ }
+ const time = new Date(savedAt).toLocaleTimeString([], {
+ hour: "numeric",
+ minute: "2-digit",
+ });
+ return `Saved ${time}`;
+}
+
+export function toErrorMessage(error: unknown, fallback: string): string {
+ if (error instanceof Error) {
+ return error.message;
+ }
+ return fallback;
+}
+
+export function normalizeDatasetRows(value: unknown): Record[] {
+ if (!Array.isArray(value)) {
+ return [];
+ }
+ return value.filter(
+ (row): row is Record =>
+ typeof row === "object" && row !== null && !Array.isArray(row),
+ );
+}
+
+export function normalizeObject(value: unknown): Record | null {
+ if (typeof value !== "object" || value === null || Array.isArray(value)) {
+ return null;
+ }
+ return value as Record;
+}
+
+export function normalizeAnalysis(value: unknown): RecipeExecutionAnalysis | null {
+ const normalized = normalizeObject(value);
+ if (!normalized) {
+ return null;
+ }
+ return normalized as RecipeExecutionAnalysis;
+}
+
+export function mapJobStatus(status: string): RecipeExecutionStatus {
+ if (status === "active") {
+ return "active";
+ }
+ if (status === "pending") {
+ return "pending";
+ }
+ if (status === "cancelling") {
+ return "cancelling";
+ }
+ if (status === "cancelled") {
+ return "cancelled";
+ }
+ if (status === "completed") {
+ return "completed";
+ }
+ if (status === "error") {
+ return "error";
+ }
+ return "running";
+}
+
+export function executionLabel(kind: "preview" | "full"): string {
+ return kind === "preview" ? "Preview" : "Full run";
+}
+
+function executionSortWeight(status: RecipeExecutionStatus): number {
+ if (status === "running" || status === "active" || status === "pending" || status === "cancelling") {
+ return 0;
+ }
+ if (status === "error" || status === "cancelled") {
+ return 2;
+ }
+ return 1;
+}
+
+export function sortExecutions(records: RecipeExecutionRecord[]): RecipeExecutionRecord[] {
+ const next = [...records];
+ next.sort((a, b) => {
+ const statusDelta = executionSortWeight(a.status) - executionSortWeight(b.status);
+ if (statusDelta !== 0) {
+ return statusDelta;
+ }
+ return b.createdAt - a.createdAt;
+ });
+ return next;
+}
+
+export function withExecutionDefaults(
+ record: RecipeExecutionRecord,
+): RecipeExecutionRecord {
+ const dataset = Array.isArray(record.dataset) ? record.dataset : [];
+ const datasetPageSize =
+ typeof record.datasetPageSize === "number" && record.datasetPageSize > 0
+ ? record.datasetPageSize
+ : DATASET_PAGE_SIZE;
+ const datasetPage =
+ typeof record.datasetPage === "number" && record.datasetPage > 0
+ ? record.datasetPage
+ : 1;
+ const datasetTotal =
+ typeof record.datasetTotal === "number" && record.datasetTotal >= 0
+ ? record.datasetTotal
+ : dataset.length;
+
+ return {
+ ...record,
+ dataset,
+ datasetTotal,
+ datasetPage,
+ datasetPageSize,
+ column_progress: record.column_progress ?? null,
+ };
+}
+
+export function delay(ms: number): Promise {
+ return new Promise((resolve) => {
+ window.setTimeout(resolve, ms);
+ });
+}
+
+export async function copyTextToClipboard(text: string): Promise {
+ try {
+ if (navigator.clipboard?.writeText) {
+ await navigator.clipboard.writeText(text);
+ return true;
+ }
+ } catch {
+ // fallthrough to legacy path
+ }
+
+ try {
+ const textarea = document.createElement("textarea");
+ textarea.value = text;
+ textarea.setAttribute("readonly", "");
+ textarea.style.position = "fixed";
+ textarea.style.top = "0";
+ textarea.style.left = "-9999px";
+ document.body.appendChild(textarea);
+ textarea.select();
+ const ok = document.execCommand("copy");
+ document.body.removeChild(textarea);
+ return ok;
+ } catch {
+ return false;
+ }
+}
diff --git a/studio/frontend/src/features/recipe-studio/hooks/use-recipe-studio-actions.ts b/studio/frontend/src/features/recipe-studio/hooks/use-recipe-studio-actions.ts
index be9162bfd0..6c53e58eb1 100644
--- a/studio/frontend/src/features/recipe-studio/hooks/use-recipe-studio-actions.ts
+++ b/studio/frontend/src/features/recipe-studio/hooks/use-recipe-studio-actions.ts
@@ -7,12 +7,26 @@ import {
getRecipeJobAnalysis,
getRecipeJobDataset,
getRecipeJobStatus,
- previewRecipe,
+ streamRecipeJobEvents,
validateRecipe,
} from "../api";
import { listRecipeExecutions, saveRecipeExecution } from "../data/executions-db";
+import {
+ buildSignature,
+ copyTextToClipboard,
+ DATASET_PAGE_SIZE,
+ delay,
+ executionLabel,
+ formatSavedLabel,
+ mapJobStatus,
+ normalizeAnalysis,
+ normalizeDatasetRows,
+ normalizeObject,
+ sortExecutions,
+ toErrorMessage,
+ withExecutionDefaults,
+} from "../executions/execution-helpers";
import type {
- RecipeExecutionAnalysis,
RecipeExecutionRecord,
RecipeExecutionStatus,
} from "../execution-types";
@@ -74,157 +88,14 @@ type UseRecipeStudioActionsResult = {
importRecipe: (value: string) => string | null;
};
-const DATASET_PAGE_SIZE = 20;
-
-function buildSignature(name: string, payload: RecipePayload): string {
- return JSON.stringify({ name, payload });
-}
-
-function formatSavedLabel(savedAt: number | null): string {
- if (!savedAt) {
- return "Not saved yet";
- }
- const time = new Date(savedAt).toLocaleTimeString([], {
- hour: "numeric",
- minute: "2-digit",
- });
- return `Saved ${time}`;
-}
-
-function toErrorMessage(error: unknown, fallback: string): string {
- if (error instanceof Error) {
- return error.message;
- }
- return fallback;
-}
-
-function normalizeDatasetRows(value: unknown): Record[] {
- if (!Array.isArray(value)) {
- return [];
- }
- return value.filter(
- (row): row is Record =>
- typeof row === "object" && row !== null && !Array.isArray(row),
- );
-}
-
-function normalizeObject(value: unknown): Record | null {
- if (typeof value !== "object" || value === null || Array.isArray(value)) {
- return null;
- }
- return value as Record;
-}
-
-function normalizeAnalysis(value: unknown): RecipeExecutionAnalysis | null {
- const normalized = normalizeObject(value);
- if (!normalized) {
- return null;
- }
- return normalized as RecipeExecutionAnalysis;
-}
-
-function mapJobStatus(status: string): RecipeExecutionStatus {
- if (status === "active") {
- return "active";
- }
- if (status === "pending") {
- return "pending";
- }
- if (status === "cancelling") {
- return "cancelling";
- }
- if (status === "cancelled") {
- return "cancelled";
- }
- if (status === "completed") {
- return "completed";
- }
- if (status === "error") {
- return "error";
- }
- return "running";
-}
-
-function executionSortWeight(status: RecipeExecutionStatus): number {
- if (status === "running" || status === "active" || status === "pending" || status === "cancelling") {
- return 0;
- }
- if (status === "error" || status === "cancelled") {
- return 2;
- }
- return 1;
-}
-
-function sortExecutions(records: RecipeExecutionRecord[]): RecipeExecutionRecord[] {
- const next = [...records];
- next.sort((a, b) => {
- const statusDelta = executionSortWeight(a.status) - executionSortWeight(b.status);
- if (statusDelta !== 0) {
- return statusDelta;
- }
- return b.createdAt - a.createdAt;
- });
- return next;
-}
-
-function withExecutionDefaults(
- record: RecipeExecutionRecord,
-): RecipeExecutionRecord {
- const dataset = Array.isArray(record.dataset) ? record.dataset : [];
- const datasetPageSize =
- typeof record.datasetPageSize === "number" && record.datasetPageSize > 0
- ? record.datasetPageSize
- : DATASET_PAGE_SIZE;
- const datasetPage =
- typeof record.datasetPage === "number" && record.datasetPage > 0
- ? record.datasetPage
- : 1;
- const datasetTotal =
- typeof record.datasetTotal === "number" && record.datasetTotal >= 0
- ? record.datasetTotal
- : dataset.length;
-
- return {
- ...record,
- dataset,
- datasetTotal,
- datasetPage,
- datasetPageSize,
- };
-}
-
-function delay(ms: number): Promise {
- return new Promise((resolve) => {
- window.setTimeout(resolve, ms);
- });
-}
-
-async function copyTextToClipboard(text: string): Promise {
- try {
- if (navigator.clipboard?.writeText) {
- await navigator.clipboard.writeText(text);
- return true;
- }
- } catch {
- // fallthrough to legacy path
- }
-
- try {
- const textarea = document.createElement("textarea");
- textarea.value = text;
- textarea.setAttribute("readonly", "");
- textarea.style.position = "fixed";
- textarea.style.top = "0";
- textarea.style.left = "-9999px";
- document.body.appendChild(textarea);
- textarea.select();
- const ok = document.execCommand("copy");
- document.body.removeChild(textarea);
- return ok;
- } catch {
- return false;
- }
-}
+type JobCompletedEventPayload = {
+ analysis?: unknown;
+ dataset?: unknown;
+ // biome-ignore lint/style/useNamingConvention: backend schema
+ artifact_path?: unknown;
+ error?: unknown;
+ type?: unknown;
+};
export function useRecipeStudioActions({
recipeId,
@@ -383,131 +254,22 @@ export function useRecipeStudioActions({
setPreviewDialogOpen(true);
}
- const runPreview = useCallback(async (): Promise => {
- const payload = readPayload();
- if (!payload) {
- setPreviewErrors(payloadResult.errors);
- toastError("Invalid recipe payload", payloadErrorMessage);
- return false;
- }
- setPreviewLoading(true);
+ const runJobExecution = useCallback(async (input: {
+ kind: "preview" | "full";
+ payload: RecipePayload;
+ rows: number;
+ }): Promise => {
+ const { kind, payload, rows } = input;
+ const setLoading = kind === "preview" ? setPreviewLoading : setFullLoading;
+ const label = executionLabel(kind);
+ setLoading(true);
const createdAt = Date.now();
const baseExecution: RecipeExecutionRecord = {
id: crypto.randomUUID(),
recipeId,
jobId: null,
- kind: "preview",
- status: "running",
- rows: previewRows,
- createdAt,
- finishedAt: null,
- recipeSignature: currentSignature,
- stage: "preview",
- current_column: null,
- progress: null,
- model_usage: null,
- lastEventId: null,
- artifact_path: null,
- dataset: [],
- datasetTotal: 0,
- datasetPage: 1,
- datasetPageSize: DATASET_PAGE_SIZE,
- analysis: null,
- processor_artifacts: null,
- error: null,
- };
- upsertExecution(baseExecution);
- onExecutionStart?.();
- setPreviewDialogOpen(false);
-
- const previewPayload = {
- ...payload,
- run: {
- ...payload.run,
- rows: previewRows,
- },
- };
-
- try {
- const validation = await validateRecipe(previewPayload);
- if (!validation.valid) {
- const errors = validation.errors.map((item) => item.message);
- const fallback = validation.raw_detail ?? "Validation failed.";
- const nextErrors = errors.length > 0 ? errors : [fallback];
- upsertExecution({
- ...baseExecution,
- status: "error",
- error: nextErrors[0],
- });
- setPreviewErrors(nextErrors);
- toastError("Validation failed", nextErrors[0]);
- return false;
- }
-
- const result = await previewRecipe(previewPayload);
- const dataset = normalizeDatasetRows(result.dataset);
- upsertExecution({
- ...baseExecution,
- status: "completed",
- finishedAt: Date.now(),
- dataset,
- datasetTotal: dataset.length,
- datasetPage: 1,
- datasetPageSize: DATASET_PAGE_SIZE,
- analysis: normalizeAnalysis(result.analysis),
- processor_artifacts: normalizeObject(result.processor_artifacts),
- error: null,
- });
- setPreviewErrors([]);
- toastSuccess(`Preview generated (${previewRows} rows).`);
- onPreviewSuccess?.();
- return true;
- } catch (error) {
- console.error("Preview failed:", error);
- const message = toErrorMessage(error, "Preview request failed.");
- upsertExecution({
- ...baseExecution,
- status: "error",
- finishedAt: Date.now(),
- error: message,
- });
- setPreviewErrors([message]);
- toastError("Preview failed", message);
- return false;
- } finally {
- setPreviewLoading(false);
- }
- }, [
- currentSignature,
- onExecutionStart,
- onPreviewSuccess,
- payloadErrorMessage,
- payloadResult.errors,
- previewRows,
- readPayload,
- recipeId,
- upsertExecution,
- ]);
-
- const runFull = useCallback(async (): Promise => {
- const payload = readPayload();
- if (!payload) {
- setPreviewErrors(payloadResult.errors);
- toastError("Invalid recipe payload", payloadErrorMessage);
- return false;
- }
-
- const requestedRows = Number(payload.run?.rows);
- const rows = Number.isFinite(requestedRows) && requestedRows > 0
- ? Math.floor(requestedRows)
- : 1000;
- const createdAt = Date.now();
- const baseExecution: RecipeExecutionRecord = {
- id: crypto.randomUUID(),
- recipeId,
- jobId: null,
- kind: "full",
+ kind,
status: "pending",
rows,
createdAt,
@@ -516,6 +278,7 @@ export function useRecipeStudioActions({
stage: "pending",
current_column: null,
progress: null,
+ column_progress: null,
model_usage: null,
lastEventId: null,
artifact_path: null,
@@ -530,65 +293,151 @@ export function useRecipeStudioActions({
upsertExecution(baseExecution);
onExecutionStart?.();
- setFullLoading(true);
+ if (kind === "preview") {
+ setPreviewDialogOpen(false);
+ }
try {
- const fullPayload = {
+ const jobPayload = {
...payload,
run: {
...payload.run,
rows,
// biome-ignore lint/style/useNamingConvention: backend schema
- execution_type: "full",
+ execution_type: kind,
},
};
- const createdJob = await createRecipeJob(fullPayload);
+ const createdJob = await createRecipeJob(jobPayload);
const jobId = createdJob.job_id;
let done = false;
let lastStatus: RecipeExecutionStatus = "pending";
+ let completedEventPayload: JobCompletedEventPayload | null = null;
let latestExecution: RecipeExecutionRecord = {
...baseExecution,
jobId,
};
upsertExecution(latestExecution);
- while (!done) {
- const status = await getRecipeJobStatus(jobId);
- const mappedStatus = mapJobStatus(status.status);
- lastStatus = mappedStatus;
+ const eventsAbortController = new AbortController();
+ void streamRecipeJobEvents({
+ jobId,
+ signal: eventsAbortController.signal,
+ onEvent: (event) => {
+ if (typeof event.id === "number") {
+ latestExecution = {
+ ...latestExecution,
+ lastEventId: event.id,
+ };
+ }
- latestExecution = {
- ...latestExecution,
- status: mappedStatus,
- rows: status.rows ?? latestExecution.rows,
- stage: status.stage ?? latestExecution.stage,
- current_column: status.current_column ?? null,
- progress: (normalizeObject(status.progress) as RecipeExecutionRecord["progress"]) ?? null,
- model_usage: normalizeObject(status.model_usage),
- artifact_path: status.artifact_path ?? latestExecution.artifact_path,
- error: status.error ?? null,
- finishedAt:
+ const eventType =
+ typeof event.payload.type === "string" ? event.payload.type : event.event;
+ if (eventType === "job.completed") {
+ lastStatus = "completed";
+ completedEventPayload = event.payload;
+ done = true;
+ latestExecution = {
+ ...latestExecution,
+ status: "completed",
+ finishedAt: Date.now(),
+ artifact_path:
+ typeof event.payload.artifact_path === "string"
+ ? event.payload.artifact_path
+ : latestExecution.artifact_path,
+ error: null,
+ };
+ upsertExecution(latestExecution);
+ return;
+ }
+
+ if (eventType === "job.error") {
+ lastStatus = "error";
+ done = true;
+ latestExecution = {
+ ...latestExecution,
+ status: "error",
+ finishedAt: Date.now(),
+ error:
+ typeof event.payload.error === "string"
+ ? event.payload.error
+ : latestExecution.error ?? `${label} failed.`,
+ };
+ upsertExecution(latestExecution);
+ return;
+ }
+
+ if (eventType === "job.cancelling") {
+ latestExecution = {
+ ...latestExecution,
+ status: "cancelling",
+ };
+ upsertExecution(latestExecution);
+ }
+ },
+ }).catch(() => {
+ // polling remains fallback source of truth
+ });
+
+ try {
+ while (!done) {
+ const status = await getRecipeJobStatus(jobId);
+ const mappedStatus = mapJobStatus(status.status);
+ lastStatus = mappedStatus;
+
+ latestExecution = {
+ ...latestExecution,
+ status: mappedStatus,
+ rows: status.rows ?? latestExecution.rows,
+ stage: status.stage ?? latestExecution.stage,
+ current_column: status.current_column ?? null,
+ progress: (normalizeObject(status.progress) as RecipeExecutionRecord["progress"]) ?? null,
+ column_progress:
+ (normalizeObject(status.column_progress) as RecipeExecutionRecord["column_progress"]) ?? null,
+ model_usage: normalizeObject(status.model_usage),
+ artifact_path: status.artifact_path ?? latestExecution.artifact_path,
+ error: status.error ?? null,
+ finishedAt:
+ mappedStatus === "completed" ||
+ mappedStatus === "error" ||
+ mappedStatus === "cancelled"
+ ? Date.now()
+ : null,
+ };
+ upsertExecution(latestExecution);
+
+ done =
mappedStatus === "completed" ||
mappedStatus === "error" ||
- mappedStatus === "cancelled"
- ? Date.now()
- : null,
- };
- upsertExecution(latestExecution);
-
- done =
- mappedStatus === "completed" ||
- mappedStatus === "error" ||
- mappedStatus === "cancelled";
- if (!done) {
- await delay(1200);
+ mappedStatus === "cancelled";
+ if (!done) {
+ await delay(1200);
+ }
}
+ } finally {
+ eventsAbortController.abort();
}
if (lastStatus === "completed") {
+ const eventAnalysis = completedEventPayload
+ ? completedEventPayload["analysis"]
+ : null;
+ const eventDataset = completedEventPayload
+ ? completedEventPayload["dataset"]
+ : null;
+ const shouldFetchPreviewDataset =
+ kind === "preview" && !Array.isArray(eventDataset);
+ const shouldFetchAnalysis =
+ !completedEventPayload ||
+ typeof eventAnalysis !== "object" ||
+ eventAnalysis === null ||
+ kind === "full";
const [analysisResult, datasetResult] = await Promise.allSettled([
- getRecipeJobAnalysis(jobId),
- getRecipeJobDataset(jobId, { limit: DATASET_PAGE_SIZE, offset: 0 }),
+ shouldFetchAnalysis
+ ? getRecipeJobAnalysis(jobId)
+ : Promise.resolve(eventAnalysis),
+ shouldFetchPreviewDataset || kind === "full"
+ ? getRecipeJobDataset(jobId, { limit: DATASET_PAGE_SIZE, offset: 0 })
+ : Promise.resolve({ dataset: eventDataset ?? [], total: rows }),
]);
const analysis =
analysisResult.status === "fulfilled"
@@ -617,7 +466,14 @@ export function useRecipeStudioActions({
error: null,
finishedAt: latestExecution.finishedAt ?? Date.now(),
});
- toastSuccess("Full run completed.");
+
+ if (kind === "preview") {
+ setPreviewErrors([]);
+ onPreviewSuccess?.();
+ toastSuccess(`Preview generated (${rows} rows).`);
+ } else {
+ toastSuccess("Full run completed.");
+ }
return true;
}
@@ -628,39 +484,110 @@ export function useRecipeStudioActions({
error: latestExecution.error ?? "Run cancelled.",
finishedAt: latestExecution.finishedAt ?? Date.now(),
});
- toastError("Full run cancelled", "The execution was cancelled.");
+ toastError(`${label} cancelled`, "The execution was cancelled.");
return false;
}
upsertExecution({
...latestExecution,
status: "error",
- error: latestExecution.error ?? "Full run failed.",
+ error: latestExecution.error ?? `${label} failed.`,
finishedAt: latestExecution.finishedAt ?? Date.now(),
});
- toastError("Full run failed", latestExecution.error ?? "Execution failed.");
+ toastError(`${label} failed`, latestExecution.error ?? "Execution failed.");
return false;
} catch (error) {
- const message = toErrorMessage(error, "Full run request failed.");
+ const message = toErrorMessage(error, `${label} request failed.`);
upsertExecution({
...baseExecution,
status: "error",
error: message,
finishedAt: Date.now(),
});
- toastError("Full run failed", message);
+ if (kind === "preview") {
+ setPreviewErrors([message]);
+ }
+ toastError(`${label} failed`, message);
return false;
} finally {
- setFullLoading(false);
+ setLoading(false);
}
}, [
currentSignature,
onExecutionStart,
+ onPreviewSuccess,
+ recipeId,
+ upsertExecution,
+ ]);
+
+ const runPreview = useCallback(async (): Promise => {
+ const payload = readPayload();
+ if (!payload) {
+ setPreviewErrors(payloadResult.errors);
+ toastError("Invalid recipe payload", payloadErrorMessage);
+ return false;
+ }
+
+ const previewPayload = {
+ ...payload,
+ run: {
+ ...payload.run,
+ rows: previewRows,
+ },
+ };
+ try {
+ const validation = await validateRecipe(previewPayload);
+ if (!validation.valid) {
+ const errors = validation.errors.map((item) => item.message);
+ const fallback = validation.raw_detail ?? "Validation failed.";
+ const nextErrors = errors.length > 0 ? errors : [fallback];
+ setPreviewErrors(nextErrors);
+ toastError("Validation failed", nextErrors[0]);
+ return false;
+ }
+ } catch (error) {
+ const message = toErrorMessage(error, "Validation failed.");
+ setPreviewErrors([message]);
+ toastError("Validation failed", message);
+ return false;
+ }
+
+ return runJobExecution({
+ kind: "preview",
+ payload,
+ rows: previewRows,
+ });
+ }, [
+ payloadErrorMessage,
+ payloadResult.errors,
+ previewRows,
+ readPayload,
+ runJobExecution,
+ ]);
+
+ const runFull = useCallback(async (): Promise => {
+ const payload = readPayload();
+ if (!payload) {
+ setPreviewErrors(payloadResult.errors);
+ toastError("Invalid recipe payload", payloadErrorMessage);
+ return false;
+ }
+
+ const requestedRows = Number(payload.run?.rows);
+ const rows = Number.isFinite(requestedRows) && requestedRows > 0
+ ? Math.floor(requestedRows)
+ : 1000;
+
+ return runJobExecution({
+ kind: "full",
+ payload,
+ rows,
+ });
+ }, [
payloadErrorMessage,
payloadResult.errors,
readPayload,
- recipeId,
- upsertExecution,
+ runJobExecution,
]);
const cancelExecution = useCallback(async (id: string): Promise => {