From 521829a7fcef5b70e161534ac6c8819b5a227dfe Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Tue, 28 Jul 2026 16:31:12 +0200 Subject: [PATCH] Add intermediate previews and graph reconnection --- README.md | 7 +- src/govoplan_dataflow/backend/executor.py | 31 ++++ src/govoplan_dataflow/backend/router.py | 5 + src/govoplan_dataflow/backend/schemas.py | 15 ++ src/govoplan_dataflow/backend/service.py | 5 + tests/test_graph_and_sql.py | 45 ++++- tests/test_service.py | 11 +- .../scripts/test-dataflow-page-structure.mjs | 10 ++ webui/src/api/dataflow.ts | 10 ++ .../src/features/dataflow/DataflowCanvas.tsx | 163 +++++++++++++++--- webui/src/features/dataflow/DataflowPage.tsx | 94 +++++++--- webui/src/features/dataflow/NodeInspector.tsx | 42 +++-- .../features/dataflow/graphInteractions.ts | 119 +++++++++++++ webui/src/styles/dataflow.css | 24 +++ 14 files changed, 520 insertions(+), 61 deletions(-) create mode 100644 webui/src/features/dataflow/graphInteractions.ts diff --git a/README.md b/README.md index bb70617..d8053ac 100644 --- a/README.md +++ b/README.md @@ -34,7 +34,8 @@ nodes by purpose: Join nodes have explicit left and right ports. Append nodes accept two or more inputs. Derived columns use a constrained operation catalogue rather than -arbitrary code. +arbitrary code. The graph editor validates manual, proximity-created, and +reconnected edges against the same port, multiplicity, and cycle rules. Governed inputs are resolved through the versioned Core capability `datasources.catalogue`; Dataflow imports neither Datasources nor Connectors and @@ -53,7 +54,9 @@ execution. Preview reads at most 250 rows per source and enforces time, intermediate-row, result-byte, graph-node, and response-row bounds. Saved previews record the pipeline revision, executor version, source fingerprints, node diagnostics, -and output summary, but not source or result rows. +and output summary, but not source or result rows. A preview may return one +explicitly requested intermediate node state under the same response bound; +that row payload remains ephemeral and is not added to run evidence. Saved revisions can also be started through the versioned `dataflow.runLifecycle` capability or the Run dialog. The first runner is diff --git a/src/govoplan_dataflow/backend/executor.py b/src/govoplan_dataflow/backend/executor.py index 45a7742..9e74653 100644 --- a/src/govoplan_dataflow/backend/executor.py +++ b/src/govoplan_dataflow/backend/executor.py @@ -13,6 +13,7 @@ from govoplan_dataflow.backend.schemas import ( DataflowDiagnostic, GraphNode, NodePreviewDiagnostic, + NodePreviewResult, PipelineGraph, PreviewColumn, ) @@ -35,6 +36,7 @@ class PipelineExecutionError(RuntimeError): source_fingerprints: tuple[dict[str, Any], ...] = (), input_row_count: int = 0, diagnostics: tuple[DataflowDiagnostic, ...] = (), + node_preview: NodePreviewResult | None = None, ) -> None: super().__init__(message) self.node_id = node_id @@ -42,6 +44,7 @@ class PipelineExecutionError(RuntimeError): self.source_fingerprints = source_fingerprints self.input_row_count = input_row_count self.diagnostics = diagnostics + self.node_preview = node_preview @dataclass(frozen=True) @@ -52,6 +55,7 @@ class PipelineExecutionResult: columns: list[PreviewColumn] diagnostics: list[DataflowDiagnostic] node_diagnostics: list[NodePreviewDiagnostic] + node_preview: NodePreviewResult | None source_fingerprints: list[dict[str, Any]] input_row_count: int @@ -74,6 +78,7 @@ def execute_preview( *, row_limit: int, source_resolver: SourceResolver | None = None, + preview_node_id: str | None = None, ) -> PipelineExecutionResult: validation = validate_graph(graph) errors = [item for item in validation if item.severity == "error"] @@ -81,6 +86,11 @@ def execute_preview( raise PipelineExecutionError(errors[0].message, node_id=errors[0].node_id) node_by_id = {node.id: node for node in graph.nodes} + if preview_node_id is not None and preview_node_id not in node_by_id: + raise PipelineExecutionError( + "The requested preview node is not part of this pipeline.", + node_id=preview_node_id, + ) inputs = graph_inputs_by_port(graph) ordered, cyclic = topological_order(graph) if cyclic: @@ -120,6 +130,7 @@ def execute_preview( source_fingerprints=tuple(source_fingerprints), input_row_count=input_row_count, diagnostics=tuple(item for item in validation if item.severity != "error"), + node_preview=_node_preview(outputs, preview_node_id, row_limit), ) try: if node.type == "source.inline": @@ -230,6 +241,7 @@ def execute_preview( source_fingerprints=tuple(source_fingerprints), input_row_count=input_row_count, diagnostics=tuple(item for item in validation if item.severity != "error"), + node_preview=_node_preview(outputs, preview_node_id, row_limit), ) from exc outputs[node.id] = output_rows @@ -255,11 +267,30 @@ def execute_preview( columns=infer_columns(all_rows), diagnostics=[item for item in validation if item.severity != "error"], node_diagnostics=node_diagnostics, + node_preview=_node_preview(outputs, preview_node_id, row_limit), source_fingerprints=source_fingerprints, input_row_count=input_row_count, ) +def _node_preview( + outputs: dict[str, list[dict[str, Any]]], + node_id: str | None, + row_limit: int, +) -> NodePreviewResult | None: + if node_id is None or node_id not in outputs: + return None + all_rows = outputs[node_id] + rows = [dict(row) for row in all_rows[:row_limit]] + return NodePreviewResult( + node_id=node_id, + columns=infer_columns(all_rows), + rows=rows, + total_rows=len(all_rows), + truncated=len(rows) < len(all_rows), + ) + + def infer_columns(rows: list[dict[str, Any]]) -> list[PreviewColumn]: names: list[str] = [] for row in rows: diff --git a/src/govoplan_dataflow/backend/router.py b/src/govoplan_dataflow/backend/router.py index 77ce783..ac37514 100644 --- a/src/govoplan_dataflow/backend/router.py +++ b/src/govoplan_dataflow/backend/router.py @@ -1096,6 +1096,11 @@ def api_preview_pipeline( "run_id": response.run_id, "status": response.status, "output_row_count": response.total_rows, + "preview_node_id": ( + response.node_preview.node_id + if response.node_preview is not None + else None + ), }, ) session.commit() diff --git a/src/govoplan_dataflow/backend/schemas.py b/src/govoplan_dataflow/backend/schemas.py index f3e4df2..e29e86f 100644 --- a/src/govoplan_dataflow/backend/schemas.py +++ b/src/govoplan_dataflow/backend/schemas.py @@ -184,6 +184,12 @@ class PipelinePreviewRequest(BaseModel): graph: PipelineGraph | None = None sql_text: str | None = Field(default=None, max_length=100_000) source_nodes: list[GraphNode] = Field(default_factory=list, max_length=20) + preview_node_id: str | None = Field( + default=None, + min_length=1, + max_length=100, + pattern=r"^[A-Za-z0-9_.:-]+$", + ) row_limit: int = Field(default=100, ge=1, le=500) @@ -203,6 +209,14 @@ class NodePreviewDiagnostic(BaseModel): messages: list[str] = Field(default_factory=list) +class NodePreviewResult(BaseModel): + node_id: str + columns: list[PreviewColumn] + rows: list[dict[str, Any]] + total_rows: int + truncated: bool + + class PipelinePreviewResponse(BaseModel): run_id: str | None pipeline_id: str | None @@ -214,6 +228,7 @@ class PipelinePreviewResponse(BaseModel): truncated: bool diagnostics: list[DataflowDiagnostic] node_diagnostics: list[NodePreviewDiagnostic] + node_preview: NodePreviewResult | None = None source_fingerprints: list[dict[str, Any]] input_row_count: int definition_hash: str diff --git a/src/govoplan_dataflow/backend/service.py b/src/govoplan_dataflow/backend/service.py index 605e293..a19ddab 100644 --- a/src/govoplan_dataflow/backend/service.py +++ b/src/govoplan_dataflow/backend/service.py @@ -567,6 +567,7 @@ def preview_pipeline( truncated=False, diagnostics=validated.diagnostics, node_diagnostics=[], + node_preview=None, source_fingerprints=[], input_row_count=0, definition_hash="", @@ -622,6 +623,7 @@ def preview_pipeline( graph, row_limit=payload.row_limit, source_resolver=resolve_source, + preview_node_id=payload.preview_node_id, ) status = "succeeded" error = None @@ -631,6 +633,7 @@ def preview_pipeline( total_rows = result.total_rows truncated = result.truncated node_diagnostics = result.node_diagnostics + node_preview = result.node_preview source_fingerprints = result.source_fingerprints input_row_count = result.input_row_count except PipelineExecutionError as exc: @@ -650,6 +653,7 @@ def preview_pipeline( total_rows = 0 truncated = False node_diagnostics = list(exc.node_diagnostics) + node_preview = exc.node_preview source_fingerprints = list(exc.source_fingerprints) input_row_count = exc.input_row_count @@ -686,6 +690,7 @@ def preview_pipeline( truncated=truncated, diagnostics=diagnostics, node_diagnostics=node_diagnostics, + node_preview=node_preview, source_fingerprints=source_fingerprints, input_row_count=input_row_count, definition_hash=graph_hash, diff --git a/tests/test_graph_and_sql.py b/tests/test_graph_and_sql.py index 58d02fc..3e6cfa2 100644 --- a/tests/test_graph_and_sql.py +++ b/tests/test_graph_and_sql.py @@ -446,6 +446,43 @@ class DataflowGraphAndSqlTests(unittest.TestCase): self.assertEqual(3, result.total_rows) self.assertTrue(result.truncated) + def test_preview_can_return_a_bounded_intermediate_node_result(self) -> None: + source = inline_source() + filter_node = GraphNode( + id="filter", + type="filter", + label="Open cases", + position=GraphPosition(x=240, y=160), + config={"column": "status", "operator": "eq", "value": "open"}, + ) + output = GraphNode( + id="output", + type="output", + label="Output", + position=GraphPosition(x=440, y=160), + config={}, + ) + graph = PipelineGraph( + nodes=[source, filter_node, output], + edges=[ + GraphEdge(id="source-filter", source="source", target="filter"), + GraphEdge(id="filter-output", source="filter", target="output"), + ], + ) + + result = execute_preview( + graph, + row_limit=1, + preview_node_id="source", + ) + + self.assertEqual(2, result.total_rows) + self.assertIsNotNone(result.node_preview) + self.assertEqual("source", result.node_preview.node_id) + self.assertEqual(3, result.node_preview.total_rows) + self.assertEqual(1, len(result.node_preview.rows)) + self.assertTrue(result.node_preview.truncated) + def test_empty_aggregate_and_null_sort_are_deterministic(self) -> None: graph, _, _ = compile_sql( """ @@ -524,7 +561,7 @@ class DataflowGraphAndSqlTests(unittest.TestCase): ) with self.assertRaises(PipelineExecutionError) as raised: - execute_preview(graph, row_limit=100) + execute_preview(graph, row_limit=100, preview_node_id="source") self.assertEqual("filter-1", raised.exception.node_id) self.assertIn("Cannot apply", str(raised.exception)) @@ -536,6 +573,12 @@ class DataflowGraphAndSqlTests(unittest.TestCase): ], ) self.assertEqual(1, len(raised.exception.source_fingerprints)) + self.assertIsNotNone(raised.exception.node_preview) + self.assertEqual("source", raised.exception.node_preview.node_id) + self.assertEqual( + [{"amount": "not-a-number"}], + raised.exception.node_preview.rows, + ) if __name__ == "__main__": diff --git a/tests/test_service.py b/tests/test_service.py index e46dc60..5d2bbf6 100644 --- a/tests/test_service.py +++ b/tests/test_service.py @@ -277,7 +277,11 @@ class DataflowServiceTests(unittest.TestCase): self.session, tenant_id="tenant-1", actor_id="user-1", - payload=PipelinePreviewRequest(pipeline_id=pipeline.id, row_limit=1), + payload=PipelinePreviewRequest( + pipeline_id=pipeline.id, + preview_node_id="source", + row_limit=1, + ), principal=principal(), ) self.session.commit() @@ -287,6 +291,11 @@ class DataflowServiceTests(unittest.TestCase): self.assertEqual([{"id": 2, "amount": 15}], response.rows) self.assertEqual(2, response.total_rows) self.assertTrue(response.truncated) + self.assertIsNotNone(response.node_preview) + self.assertEqual("source", response.node_preview.node_id) + self.assertEqual([{"id": 1, "amount": 5}], response.node_preview.rows) + self.assertEqual(3, response.node_preview.total_rows) + self.assertTrue(response.node_preview.truncated) self.assertEqual(2, run.output_row_count) self.assertEqual(3, run.input_row_count) self.assertEqual(1, len(run.source_fingerprints)) diff --git a/webui/scripts/test-dataflow-page-structure.mjs b/webui/scripts/test-dataflow-page-structure.mjs index 503527f..2becaa8 100644 --- a/webui/scripts/test-dataflow-page-structure.mjs +++ b/webui/scripts/test-dataflow-page-structure.mjs @@ -2,6 +2,7 @@ import { readFileSync } from "node:fs"; const page = readFileSync(new URL("../src/features/dataflow/DataflowPage.tsx", import.meta.url), "utf8"); const canvas = readFileSync(new URL("../src/features/dataflow/DataflowCanvas.tsx", import.meta.url), "utf8"); +const interactions = readFileSync(new URL("../src/features/dataflow/graphInteractions.ts", import.meta.url), "utf8"); const css = readFileSync(new URL("../src/styles/dataflow.css", import.meta.url), "utf8"); const moduleEntry = readFileSync(new URL("../src/module.ts", import.meta.url), "utf8"); @@ -16,6 +17,15 @@ const checks = [ [canvas.includes("connectionRadius={32}"), "visible connection drop radius"], [canvas.includes("initialWidth: 180"), "known initial node width"], [canvas.includes('change.type !== "dimensions"'), "measurement-only change filtering"], + [canvas.includes("onReconnect={onReconnect}"), "edge reconnection"], + [canvas.includes("onReconnectEnd="), "edge deletion on dropped reconnection"], + [canvas.includes("selectedEdgeId"), "persistent edge selection"], + [ + canvas.includes("closestProximityConnection") + && interactions.includes("definitionConnectionError"), + "validated proximity connection" + ], + [css.includes(".dataflow-edge-proximity"), "proximity connection indicator"], [ moduleEntry.indexOf("@xyflow/react/dist/style.css") < moduleEntry.indexOf("./styles/dataflow.css"), diff --git a/webui/src/api/dataflow.ts b/webui/src/api/dataflow.ts index 36c7542..b5d1ba7 100644 --- a/webui/src/api/dataflow.ts +++ b/webui/src/api/dataflow.ts @@ -125,6 +125,14 @@ export type NodePreviewDiagnostic = { messages: string[]; }; +export type NodePreviewResult = { + node_id: string; + columns: PreviewColumn[]; + rows: Record[]; + total_rows: number; + truncated: boolean; +}; + export type PipelinePreview = { run_id?: string | null; pipeline_id?: string | null; @@ -136,6 +144,7 @@ export type PipelinePreview = { truncated: boolean; diagnostics: DataflowDiagnostic[]; node_diagnostics: NodePreviewDiagnostic[]; + node_preview?: NodePreviewResult | null; source_fingerprints: Record[]; input_row_count: number; definition_hash: string; @@ -451,6 +460,7 @@ export function previewDataflowPipeline( graph?: PipelineGraph; sql_text?: string; source_nodes?: PipelineGraphNode[]; + preview_node_id?: string; row_limit?: number; } ): Promise { diff --git a/webui/src/features/dataflow/DataflowCanvas.tsx b/webui/src/features/dataflow/DataflowCanvas.tsx index 2887d1b..f53927d 100644 --- a/webui/src/features/dataflow/DataflowCanvas.tsx +++ b/webui/src/features/dataflow/DataflowCanvas.tsx @@ -1,4 +1,4 @@ -import { useMemo, useState, type DragEvent } from "react"; +import { useMemo, useRef, useState, type DragEvent } from "react"; import { addEdge, applyEdgeChanges, @@ -9,6 +9,7 @@ import { Controls, MiniMap, ReactFlow, + reconnectEdge, type Connection, type Edge, type ReactFlowInstance @@ -22,6 +23,11 @@ import type { PipelineGraphNode } from "../../api/dataflow"; import DataflowNode, { type DataflowFlowNode } from "./DataflowNode"; +import { + closestProximityConnection, + graphEdgesFromFlow, + proximityPreviewEdge +} from "./graphInteractions"; import { FALLBACK_NODE_LIBRARY, newNode } from "./model"; const nodeTypes = { dataflow: DataflowNode }; @@ -48,6 +54,10 @@ export default function DataflowCanvas({ onSelectNode }: DataflowCanvasProps) { const [instance, setInstance] = useState | null>(null); + const [proximityEdge, setProximityEdge] = useState(null); + const [selectedEdgeId, setSelectedEdgeId] = useState(null); + const reconnectSuccessful = useRef(true); + const reconnectingEdgeId = useRef(null); const errorNodeIds = useMemo( () => new Set(diagnostics.filter((item) => item.severity === "error" && item.node_id).map((item) => item.node_id)), [diagnostics] @@ -92,9 +102,14 @@ export default function DataflowCanvas({ sourceHandle: edge.source_port ?? "output", targetHandle: edge.target_port ?? "input", type: "smoothstep", - className: "dataflow-edge" + className: "dataflow-edge", + selected: edge.id === selectedEdgeId })), - [graph.edges] + [graph.edges, selectedEdgeId] + ); + const displayedEdges = useMemo( + () => proximityEdge ? [...edges, proximityEdge] : edges, + [edges, proximityEdge] ); const updateNodes = (nextNodes: DataflowFlowNode[]) => { @@ -116,28 +131,32 @@ export default function DataflowCanvas({ const updateEdges = (nextEdges: Edge[]) => { onGraphChange({ ...graph, - edges: nextEdges.map((edge) => ({ - id: edge.id, - source: edge.source, - target: edge.target, - source_port: edge.sourceHandle ?? "output", - target_port: edge.targetHandle ?? "input" - })) + edges: graphEdgesFromFlow(nextEdges) }); }; + const connectionError = ( + connection: Connection | Edge, + ignoredEdgeId: string | null = reconnectingEdgeId.current + ): string | null => definitionConnectionError( + ignoredEdgeId + ? { + ...graph, + edges: graph.edges.filter((edge) => edge.id !== ignoredEdgeId) + } + : graph, + nodeLibrary, + { + source: connection.source, + target: connection.target, + sourcePort: connection.sourceHandle, + targetPort: connection.targetHandle + } + ); + const isValidConnection = (connection: Connection | Edge): boolean => { if (readOnly || !connection.source || !connection.target) return false; - return definitionConnectionError( - graph, - nodeLibrary, - { - source: connection.source, - target: connection.target, - sourcePort: connection.sourceHandle, - targetPort: connection.targetHandle - } - ) === null; + return connectionError(connection) === null; }; const onConnect = (connection: Connection) => { @@ -153,6 +172,26 @@ export default function DataflowCanvas({ updateEdges(next); }; + const onReconnect = (oldEdge: Edge, connection: Connection) => { + if (readOnly || connectionError(connection, oldEdge.id)) return; + reconnectSuccessful.current = true; + updateEdges(reconnectEdge( + oldEdge, + connection, + edges, + { shouldReplaceId: false } + )); + }; + + const proximityConnection = (draggedNode: DataflowFlowNode) => ( + closestProximityConnection( + draggedNode, + nodes, + graph, + nodeLibrary + ) + ); + const onDrop = (event: DragEvent) => { event.preventDefault(); if (readOnly || !instance) return; @@ -177,7 +216,7 @@ export default function DataflowCanvas({ > nodes={nodes} - edges={edges} + edges={displayedEdges} nodeTypes={nodeTypes} onInit={setInstance} onNodesChange={(changes) => { @@ -189,14 +228,90 @@ export default function DataflowCanvas({ }} onEdgesChange={(changes) => { if (readOnly) return; - updateEdges(applyEdgeChanges(changes, edges)); + const selectedChange = changes.find( + (change) => change.type === "select" && change.selected + ); + if (selectedChange?.type === "select") { + setSelectedEdgeId(selectedChange.id); + } else if (changes.some( + (change) => ( + change.type === "select" + && change.id === selectedEdgeId + && !change.selected + ) + )) { + setSelectedEdgeId(null); + } + const graphChanges = changes.filter((change) => change.type !== "select"); + if (graphChanges.length) { + if (graphChanges.some( + (change) => change.type === "remove" && change.id === selectedEdgeId + )) { + setSelectedEdgeId(null); + } + updateEdges(applyEdgeChanges(graphChanges, edges)); + } }} onConnect={onConnect} + onReconnect={onReconnect} + onReconnectStart={(_event, edge) => { + reconnectSuccessful.current = false; + reconnectingEdgeId.current = edge.id; + setProximityEdge(null); + }} + onReconnectEnd={(_event, edge) => { + if (!reconnectSuccessful.current && !readOnly) { + updateEdges(edges.filter((candidate) => candidate.id !== edge.id)); + setSelectedEdgeId(null); + } + reconnectSuccessful.current = true; + reconnectingEdgeId.current = null; + }} + onNodeDrag={(_event, node) => { + if (readOnly) return; + const connection = proximityConnection(node); + setProximityEdge(connection ? proximityPreviewEdge(connection) : null); + }} + onNodeDragStop={(_event, node) => { + setProximityEdge(null); + if (readOnly) return; + const connection = proximityConnection(node); + if (!connection) return; + const nextEdges = addEdge( + { + ...connection, + id: `edge-${crypto.randomUUID()}`, + type: "smoothstep", + className: "dataflow-edge" + }, + edges + ); + onGraphChange({ + ...graph, + nodes: graph.nodes.map((graphNode) => ( + graphNode.id === node.id + ? { ...graphNode, position: node.position } + : graphNode + )), + edges: graphEdgesFromFlow(nextEdges) + }); + }} isValidConnection={isValidConnection} - onNodeClick={(_event, node) => onSelectNode(node.id)} - onPaneClick={() => onSelectNode(null)} + onNodeClick={(_event, node) => { + setSelectedEdgeId(null); + onSelectNode(node.id); + }} + onEdgeClick={(_event, edge) => { + setSelectedEdgeId(edge.id); + onSelectNode(null); + }} + onPaneClick={() => { + setSelectedEdgeId(null); + onSelectNode(null); + }} nodesDraggable={!readOnly} nodesConnectable={!readOnly} + edgesReconnectable={!readOnly} elementsSelectable deleteKeyCode={readOnly ? null : ["Backspace", "Delete"]} connectionLineType={ConnectionLineType.SmoothStep} diff --git a/webui/src/features/dataflow/DataflowPage.tsx b/webui/src/features/dataflow/DataflowPage.tsx index b378dfc..61451fb 100644 --- a/webui/src/features/dataflow/DataflowPage.tsx +++ b/webui/src/features/dataflow/DataflowPage.tsx @@ -348,6 +348,7 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings updateDraft({ graph }); setDiagnostics([]); setNodeDiagnostics([]); + setPreview(null); }; const validate = async () => { @@ -425,27 +426,45 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings await applySql(true); }; - const runPreview = async () => { - if (!draft || !canRun) return; + const runPreview = async (previewNodeId?: string) => { + if (!draft || !canPreview) return; + const targetNodeId = previewNodeId + ?? draft.graph.nodes.find((node) => node.type === "output")?.id; setWorking(true); setError(""); setSuccess(""); try { const response = await previewDataflowPipeline(settings, !dirty && draft.id && draft.currentRevision - ? { pipeline_id: draft.id, revision: draft.currentRevision, row_limit: 100 } + ? { + pipeline_id: draft.id, + revision: draft.currentRevision, + preview_node_id: targetNodeId, + row_limit: 100 + } : { graph: draft.graph, sql_text: draft.editorMode === "sql" ? draft.sqlText : undefined, source_nodes: sourceNodes(draft.graph), + preview_node_id: targetNodeId, row_limit: 100 }); setPreview(response); setDiagnostics(response.diagnostics); setNodeDiagnostics(response.node_diagnostics); setResultOpen(true); - setResultTab(response.status === "succeeded" ? "preview" : "diagnostics"); - if (response.status === "succeeded") { - setSuccess(`Preview produced ${response.total_rows} row${response.total_rows === 1 ? "" : "s"}.`); + setResultTab( + response.status === "succeeded" || response.node_preview + ? "preview" + : "diagnostics" + ); + if (response.status === "succeeded" && response.node_preview) { + const node = draft.graph.nodes.find( + (candidate) => candidate.id === response.node_preview?.node_id + ); + setSuccess( + `${node?.label ?? "Preview"} produced ${response.node_preview.total_rows} ` + + `row${response.node_preview.total_rows === 1 ? "" : "s"}.` + ); } } catch (previewError) { setError(apiErrorMessage(previewError)); @@ -754,7 +773,10 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings sources={sources} sourceCatalogueAvailable={sourceCatalogueAvailable} readOnly={!canEdit} + canPreview={canPreview} + previewing={working} onChange={(node: PipelineGraphNode) => updateGraph(updateGraphNode(draft.graph, node))} + onPreview={(nodeId) => void runPreview(nodeId)} onDelete={(nodeId) => { updateGraph({ ...draft.graph, @@ -770,8 +792,11 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings tab={resultTab} onTabChange={setResultTab} preview={preview} + nodes={draft.graph.nodes} + working={working} diagnostics={diagnostics} nodeDiagnostics={nodeDiagnostics} + onPreviewNodeChange={(nodeId) => void runPreview(nodeId)} onClose={() => setResultOpen(false)} /> ) : null} @@ -1846,29 +1871,54 @@ function ResultPanel({ tab, onTabChange, preview, + nodes, + working, diagnostics, nodeDiagnostics, + onPreviewNodeChange, onClose }: { tab: ResultTab; onTabChange: (tab: ResultTab) => void; preview: PipelinePreview | null; + nodes: PipelineGraphNode[]; + working: boolean; diagnostics: DataflowDiagnostic[]; nodeDiagnostics: NodePreviewDiagnostic[]; + onPreviewNodeChange: (nodeId: string) => void; onClose: () => void; }) { + const outputNodeId = nodes.find((node) => node.type === "output")?.id ?? ""; + const previewNodeId = preview?.node_preview?.node_id ?? outputNodeId; + const previewRows = preview?.node_preview?.total_rows ?? preview?.total_rows; return (
- - ariaLabel="Result view" - options={[ - { id: "preview", label: `Preview${preview ? ` (${preview.total_rows})` : ""}` }, - { id: "diagnostics", label: `Diagnostics (${diagnostics.length + nodeDiagnostics.length})` } - ]} - value={tab} - onChange={onTabChange} - /> +
+ + ariaLabel="Result view" + options={[ + { id: "preview", label: `Preview${previewRows !== undefined ? ` (${previewRows})` : ""}` }, + { id: "diagnostics", label: `Diagnostics (${diagnostics.length + nodeDiagnostics.length})` } + ]} + value={tab} + onChange={onTabChange} + /> + {tab === "preview" && nodes.length ? ( + + ) : null} +
{tab === "preview" ? ( @@ -1882,13 +1932,16 @@ function ResultPanel({ function PreviewTable({ preview }: { preview: PipelinePreview | null }) { if (!preview) return
No preview has been run.
; - if (preview.status === "failed") return
Preview failed.
; + if (preview.status === "failed" && !preview.node_preview) { + return
Preview failed.
; + } + const result = preview.node_preview ?? preview; return (
- {preview.columns.map((column) => ( + {result.columns.map((column) => ( - {preview.rows.map((row, rowIndex) => ( + {result.rows.map((row, rowIndex) => ( - {preview.columns.map((column) => ( + {result.columns.map((column) => ( ))} @@ -1908,8 +1961,9 @@ function PreviewTable({ preview }: { preview: PipelinePreview | null }) {
{column.name} {column.type}{column.nullable ? " ยท nullable" : ""} @@ -1897,9 +1950,9 @@ function PreviewTable({ preview }: { preview: PipelinePreview | null }) {
{formatCell(row[column.name])}
- Showing {preview.rows.length} of {preview.total_rows} output rows + Showing {result.rows.length} of {result.total_rows} rows + {preview.status === "failed" ? Pipeline failed after this stage : null} {preview.input_row_count} input rows {preview.source_fingerprints.length} source diff --git a/webui/src/features/dataflow/NodeInspector.tsx b/webui/src/features/dataflow/NodeInspector.tsx index df33c8a..f84a8f0 100644 --- a/webui/src/features/dataflow/NodeInspector.tsx +++ b/webui/src/features/dataflow/NodeInspector.tsx @@ -1,6 +1,10 @@ import { useEffect, useState } from "react"; -import { Trash2 } from "lucide-react"; -import { Button, DismissibleAlert, FormField } from "@govoplan/core-webui"; +import { Play, Trash2 } from "lucide-react"; +import { + DismissibleAlert, + FormField, + IconButton +} from "@govoplan/core-webui"; import type { NodeTypeDefinition, PipelineGraphNode, @@ -13,8 +17,11 @@ type NodeInspectorProps = { sources: TabularSource[]; sourceCatalogueAvailable: boolean; readOnly: boolean; + canPreview: boolean; + previewing: boolean; onChange: (node: PipelineGraphNode) => void; onDelete: (nodeId: string) => void; + onPreview: (nodeId: string) => void; }; export default function NodeInspector({ @@ -23,8 +30,11 @@ export default function NodeInspector({ sources, sourceCatalogueAvailable, readOnly, + canPreview, + previewing, onChange, - onDelete + onDelete, + onPreview }: NodeInspectorProps) { const [rowsText, setRowsText] = useState(""); const [aggregateText, setAggregateText] = useState(""); @@ -94,16 +104,22 @@ export default function NodeInspector({ Inspector {definition?.label ?? node.type} - + + } + variant="ghost" + onClick={() => onPreview(node.id)} + disabled={!canPreview || previewing} + /> + } + variant="danger" + onClick={() => onDelete(node.id)} + disabled={readOnly} + /> +
{localError ? ( diff --git a/webui/src/features/dataflow/graphInteractions.ts b/webui/src/features/dataflow/graphInteractions.ts new file mode 100644 index 0000000..5b91261 --- /dev/null +++ b/webui/src/features/dataflow/graphInteractions.ts @@ -0,0 +1,119 @@ +import { definitionConnectionError } from "@govoplan/core-webui/definition-graph"; +import type { Connection, Edge } from "@xyflow/react"; +import type { + NodeTypeDefinition, + PipelineGraph +} from "../../api/dataflow"; +import type { DataflowFlowNode } from "./DataflowNode"; + +const PROXIMITY_GAP = 72; +const PROXIMITY_EDGE_ID = "dataflow-proximity-preview"; + +export function graphEdgesFromFlow(edges: Edge[]): PipelineGraph["edges"] { + return edges + .filter((edge) => edge.id !== PROXIMITY_EDGE_ID) + .map((edge) => ({ + id: edge.id, + source: edge.source, + target: edge.target, + source_port: edge.sourceHandle ?? "output", + target_port: edge.targetHandle ?? "input" + })); +} + +export function proximityPreviewEdge(connection: Connection): Edge { + return { + ...connection, + id: PROXIMITY_EDGE_ID, + type: "smoothstep", + className: "dataflow-edge dataflow-edge-proximity", + animated: true, + selectable: false, + reconnectable: false + }; +} + +export function closestProximityConnection( + draggedNode: DataflowFlowNode, + flowNodes: DataflowFlowNode[], + graph: PipelineGraph, + nodeLibrary: NodeTypeDefinition[] +): Connection | null { + const currentNodes = flowNodes.map((node) => ( + node.id === draggedNode.id ? draggedNode : node + )); + const neighbors = currentNodes + .filter((node) => node.id !== draggedNode.id) + .map((node) => ({ + node, + distance: nodeGap(draggedNode, node) + })) + .filter((candidate) => candidate.distance <= PROXIMITY_GAP) + .sort((left, right) => left.distance - right.distance); + + for (const { node: neighbor } of neighbors) { + const draggedCenter = nodeCenter(draggedNode); + const neighborCenter = nodeCenter(neighbor); + const orientations: Array<[DataflowFlowNode, DataflowFlowNode]> = ( + draggedCenter.x <= neighborCenter.x + ? [[draggedNode, neighbor], [neighbor, draggedNode]] + : [[neighbor, draggedNode], [draggedNode, neighbor]] + ); + for (const [source, target] of orientations) { + for (const sourcePort of source.data.definition.output_ports) { + for (const targetPort of target.data.definition.input_ports) { + const connection: Connection = { + source: source.id, + target: target.id, + sourceHandle: sourcePort.id, + targetHandle: targetPort.id + }; + if (definitionConnectionError( + graph, + nodeLibrary, + { + source: connection.source, + target: connection.target, + sourcePort: connection.sourceHandle, + targetPort: connection.targetHandle + } + ) === null) { + return connection; + } + } + } + } + } + return null; +} + +function nodeGap(left: DataflowFlowNode, right: DataflowFlowNode): number { + const leftCenter = nodeCenter(left); + const rightCenter = nodeCenter(right); + const horizontal = Math.max( + 0, + Math.abs(leftCenter.x - rightCenter.x) + - (nodeWidth(left) + nodeWidth(right)) / 2 + ); + const vertical = Math.max( + 0, + Math.abs(leftCenter.y - rightCenter.y) + - (nodeHeight(left) + nodeHeight(right)) / 2 + ); + return Math.hypot(horizontal, vertical); +} + +function nodeCenter(node: DataflowFlowNode): { x: number; y: number } { + return { + x: node.position.x + nodeWidth(node) / 2, + y: node.position.y + nodeHeight(node) / 2 + }; +} + +function nodeWidth(node: DataflowFlowNode): number { + return node.measured?.width ?? node.initialWidth ?? 180; +} + +function nodeHeight(node: DataflowFlowNode): number { + return node.measured?.height ?? node.initialHeight ?? 54; +} diff --git a/webui/src/styles/dataflow.css b/webui/src/styles/dataflow.css index c8e5c9f..f09baca 100644 --- a/webui/src/styles/dataflow.css +++ b/webui/src/styles/dataflow.css @@ -537,6 +537,12 @@ stroke-width: 3; } +.dataflow-canvas .react-flow__edge.dataflow-edge-proximity .react-flow__edge-path { + stroke: var(--accent); + stroke-width: 3; + stroke-dasharray: 7 6; +} + .dataflow-canvas .react-flow__connection-path { stroke: var(--accent); stroke-width: 3; @@ -898,12 +904,30 @@ .dataflow-results-toolbar { min-height: 42px; padding: 5px 9px; + flex-wrap: wrap; } .dataflow-results-toolbar .btn { min-height: 30px; } +.dataflow-results-view-controls, +.dataflow-inspector-actions { + display: flex; + align-items: center; + gap: 7px; + min-width: 0; +} + +.dataflow-results-view-controls { + flex-wrap: wrap; +} + +.dataflow-results-view-controls select { + width: clamp(150px, 22vw, 280px); + min-height: 30px; +} + .dataflow-preview-table-wrap, .dataflow-diagnostics-list { min-height: 0;