Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4d570b5e9e | ||
|
|
a6c5bab3a4 | ||
|
|
4175262b8b | ||
|
|
618f10fe89 | ||
|
|
f222be63b2 |
@@ -213,3 +213,21 @@ npm run test:structure
|
||||
|
||||
The implementation epic is
|
||||
[`govoplan-dataflow#1`](https://git.add-ideas.de/GovOPlaN/govoplan-dataflow/issues/1).
|
||||
|
||||
## Git-source WebUI package
|
||||
|
||||
The repository root exposes `@govoplan/dataflow-webui` for Git-tagged release
|
||||
dependencies. It mirrors the owning `webui/package.json` version, public
|
||||
TypeScript/CSS exports and peer requirements, with entry paths under
|
||||
`webui/src`. Consumers provide the shared Core/React peers; the facade runs no
|
||||
development or install scripts. The source archive contains `webui/src`, this
|
||||
README and any repository license file. Run module development checks from `webui/`; Python
|
||||
installation remains governed by `pyproject.toml`.
|
||||
|
||||
Das Repository stellt `@govoplan/dataflow-webui` am Wurzelpfad für versionierte
|
||||
Git-Abhängigkeiten bereit. Version, öffentliche TypeScript-/CSS-Exporte und
|
||||
Peer-Anforderungen entsprechen `webui/package.json`; die Einstiegspfade liegen
|
||||
unter `webui/src`. Gemeinsame Core-/React-Peers stellt die einbindende Anwendung
|
||||
bereit. Die Fassade führt keine Entwicklungs- oder Installationsskripte aus.
|
||||
Entwicklungsprüfungen bleiben in `webui/`, die Python-Installation weiterhin in
|
||||
`pyproject.toml` definiert.
|
||||
|
||||
@@ -0,0 +1,35 @@
|
||||
{
|
||||
"name": "@govoplan/dataflow-webui",
|
||||
"version": "0.1.25",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"main": "webui/src/index.ts",
|
||||
"module": "webui/src/index.ts",
|
||||
"types": "webui/src/index.ts",
|
||||
"exports": {
|
||||
".": {
|
||||
"types": "./webui/src/index.ts",
|
||||
"import": "./webui/src/index.ts"
|
||||
},
|
||||
"./styles/dataflow.css": "./webui/src/styles/dataflow.css"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"@govoplan/core-webui": "^0.1.46",
|
||||
"@xyflow/react": "^12.11.2",
|
||||
"lucide-react": "^1.23.0",
|
||||
"react": ">=19.2.7 <20",
|
||||
"react-dom": ">=19.2.7 <20",
|
||||
"react-router": ">=8.3.0 <9",
|
||||
"typescript": "^5.7.2"
|
||||
},
|
||||
"peerDependenciesMeta": {
|
||||
"@govoplan/core-webui": {
|
||||
"optional": true
|
||||
}
|
||||
},
|
||||
"files": [
|
||||
"webui/src",
|
||||
"README.md",
|
||||
"LICENSE"
|
||||
]
|
||||
}
|
||||
+2
-2
@@ -4,14 +4,14 @@ build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "govoplan-dataflow"
|
||||
version = "0.1.23"
|
||||
version = "0.1.25"
|
||||
description = "Governed graphical and SQL data pipelines for GovOPlaN."
|
||||
readme = "README.md"
|
||||
requires-python = ">=3.12"
|
||||
license = "AGPL-3.0-or-later"
|
||||
authors = [{ name = "GovOPlaN" }]
|
||||
dependencies = [
|
||||
"govoplan-core>=0.1.18",
|
||||
"govoplan-core>=0.1.46",
|
||||
"sqlglot>=30.14,<31",
|
||||
]
|
||||
|
||||
|
||||
@@ -1,3 +1,3 @@
|
||||
from __future__ import annotations
|
||||
|
||||
__version__ = "0.1.23"
|
||||
__version__ = "0.1.25"
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import dataclass, field
|
||||
import math
|
||||
from typing import Any, Mapping, Protocol, runtime_checkable
|
||||
|
||||
from govoplan_dataflow.backend.batches import TypedBatch
|
||||
@@ -33,8 +34,8 @@ class ExecutionBudget:
|
||||
raise ValueError("Execution output row limit must be positive.")
|
||||
if self.max_batch_bytes < 1:
|
||||
raise ValueError("Execution byte limit must be positive.")
|
||||
if self.max_wall_seconds <= 0:
|
||||
raise ValueError("Execution time limit must be positive.")
|
||||
if not math.isfinite(self.max_wall_seconds) or self.max_wall_seconds <= 0:
|
||||
raise ValueError("Execution time limit must be finite and positive.")
|
||||
if self.max_memory_bytes < 64 * 1024 * 1024:
|
||||
raise ValueError("Execution memory limit must be at least 64 MiB.")
|
||||
if self.max_concurrency < 1:
|
||||
@@ -106,11 +107,19 @@ class BackendExecutionError(RuntimeError):
|
||||
code: str = "backend.execution",
|
||||
node_id: str | None = None,
|
||||
diagnostics: tuple[DataflowDiagnostic, ...] = (),
|
||||
node_diagnostics: tuple[NodePreviewDiagnostic, ...] = (),
|
||||
source_fingerprints: tuple[dict[str, Any], ...] = (),
|
||||
input_row_count: int = 0,
|
||||
node_preview: NodePreviewResult | None = None,
|
||||
) -> None:
|
||||
super().__init__(message)
|
||||
self.code = code
|
||||
self.node_id = node_id
|
||||
self.diagnostics = diagnostics
|
||||
self.node_diagnostics = node_diagnostics
|
||||
self.source_fingerprints = source_fingerprints
|
||||
self.input_row_count = input_row_count
|
||||
self.node_preview = node_preview
|
||||
|
||||
|
||||
@runtime_checkable
|
||||
|
||||
@@ -1,10 +1,25 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import asdict
|
||||
import hashlib
|
||||
import json
|
||||
import math
|
||||
|
||||
from govoplan_core.security.bounded_process import (
|
||||
ProcessBudgetError,
|
||||
ProcessLimits,
|
||||
run_bounded_operation,
|
||||
)
|
||||
from govoplan_core.security.worker_payload import (
|
||||
decode_worker_payload,
|
||||
encode_worker_payload,
|
||||
)
|
||||
from govoplan_dataflow.backend.backends.base import (
|
||||
BackendExecutionError,
|
||||
BackendExecutionRequest,
|
||||
BackendExecutionResult,
|
||||
BackendSource,
|
||||
ExecutionBudget,
|
||||
canonical_result_schema,
|
||||
)
|
||||
from govoplan_dataflow.backend.batches import TypedBatch
|
||||
@@ -14,8 +29,17 @@ from govoplan_dataflow.backend.executor import (
|
||||
ResolvedSource,
|
||||
execute_preview,
|
||||
)
|
||||
from govoplan_dataflow.backend.ir import IrExecutionResult, ir_to_graph
|
||||
from govoplan_dataflow.backend.schemas import GraphNode
|
||||
from govoplan_dataflow.backend.ir import IrExecutionResult, IrSchema, TypedGraphIr, ir_to_graph
|
||||
from govoplan_dataflow.backend.planner import ExecutionPlan
|
||||
from govoplan_dataflow.backend.schemas import (
|
||||
DataflowDiagnostic,
|
||||
GraphNode,
|
||||
NodePreviewDiagnostic,
|
||||
NodePreviewResult,
|
||||
)
|
||||
|
||||
|
||||
_TRANSPORT_BYTES = 32 * 1024 * 1024
|
||||
|
||||
|
||||
class ReferenceExecutionBackend:
|
||||
@@ -32,6 +56,52 @@ class ReferenceExecutionBackend:
|
||||
self,
|
||||
request: BackendExecutionRequest,
|
||||
) -> BackendExecutionResult:
|
||||
_validate_source_batches(request)
|
||||
try:
|
||||
limits = ProcessLimits(
|
||||
wall_seconds=request.budget.max_wall_seconds,
|
||||
cpu_seconds=max(1, math.ceil(request.budget.max_wall_seconds)),
|
||||
memory_bytes=request.budget.max_memory_bytes,
|
||||
input_bytes=_TRANSPORT_BYTES,
|
||||
output_bytes=_TRANSPORT_BYTES,
|
||||
)
|
||||
payload = encode_worker_payload(_request_payload(request), max_bytes=_TRANSPORT_BYTES)
|
||||
response = decode_worker_payload(
|
||||
run_bounded_operation(_execute_reference_worker, payload, limits=limits),
|
||||
max_bytes=_TRANSPORT_BYTES,
|
||||
)
|
||||
except ProcessBudgetError as exc:
|
||||
raise BackendExecutionError(str(exc), code=f"backend.process.{exc.code}") from exc
|
||||
except ValueError as exc:
|
||||
raise BackendExecutionError(
|
||||
"Reference execution exceeds its supported process/transport budget.",
|
||||
code="backend.budget",
|
||||
) from exc
|
||||
if "error" in response:
|
||||
error = response["error"]
|
||||
raise BackendExecutionError(
|
||||
error["message"], code=error["code"], node_id=error["node_id"],
|
||||
diagnostics=tuple(DataflowDiagnostic.model_validate(item) for item in error["diagnostics"]),
|
||||
node_diagnostics=tuple(NodePreviewDiagnostic.model_validate(item) for item in error["node_diagnostics"]),
|
||||
source_fingerprints=error["source_fingerprints"],
|
||||
input_row_count=error["input_row_count"],
|
||||
node_preview=NodePreviewResult.model_validate(error["node_preview"]) if error["node_preview"] else None,
|
||||
)
|
||||
batch = _batch_from_payload(response["batch"])
|
||||
batch.ensure_within(
|
||||
max_rows=request.budget.max_output_rows,
|
||||
max_bytes=request.budget.max_batch_bytes,
|
||||
)
|
||||
return BackendExecutionResult(
|
||||
contract=IrExecutionResult.model_validate(response["contract"]),
|
||||
batch=batch,
|
||||
node_diagnostics=tuple(NodePreviewDiagnostic.model_validate(item) for item in response["node_diagnostics"]),
|
||||
node_preview=NodePreviewResult.model_validate(response["node_preview"]) if response["node_preview"] else None,
|
||||
metadata=response["metadata"],
|
||||
)
|
||||
|
||||
def _execute_in_process(self, request: BackendExecutionRequest) -> BackendExecutionResult:
|
||||
"""Pure reference evaluation, called only inside the disposable worker."""
|
||||
_validate_source_batches(request)
|
||||
try:
|
||||
result = execute_preview(
|
||||
@@ -50,6 +120,10 @@ class ReferenceExecutionBackend:
|
||||
code="backend.reference",
|
||||
node_id=exc.node_id,
|
||||
diagnostics=tuple(exc.diagnostics),
|
||||
node_diagnostics=tuple(exc.node_diagnostics),
|
||||
source_fingerprints=tuple(exc.source_fingerprints),
|
||||
input_row_count=exc.input_row_count,
|
||||
node_preview=exc.node_preview,
|
||||
) from exc
|
||||
observed_batch = TypedBatch.from_rows(result.rows)
|
||||
batch = TypedBatch.from_rows(
|
||||
@@ -97,6 +171,90 @@ class ReferenceExecutionBackend:
|
||||
)
|
||||
|
||||
|
||||
def _batch_payload(batch: TypedBatch) -> dict:
|
||||
return {
|
||||
"schema": batch.schema.model_dump(mode="python"),
|
||||
"columns": dict(batch.columns),
|
||||
"row_count": batch.row_count,
|
||||
"byte_count": batch.byte_count,
|
||||
}
|
||||
|
||||
|
||||
def _batch_from_payload(value: dict) -> TypedBatch:
|
||||
return TypedBatch(
|
||||
schema=IrSchema.model_validate(value["schema"]), columns=value["columns"],
|
||||
row_count=value["row_count"], byte_count=value["byte_count"],
|
||||
)
|
||||
|
||||
|
||||
def _request_payload(request: BackendExecutionRequest) -> dict:
|
||||
plan = request.plan
|
||||
return {
|
||||
"plan": {
|
||||
"graph": plan.graph.model_dump(mode="python"),
|
||||
"ordered_node_ids": plan.ordered_node_ids,
|
||||
"diagnostics": tuple(item.model_dump(mode="python") for item in plan.diagnostics),
|
||||
"generated_sql": plan.generated_sql,
|
||||
"sql_diagnostics": tuple(item.model_dump(mode="python") for item in plan.sql_diagnostics),
|
||||
"semantic_hash": plan.semantic_hash,
|
||||
},
|
||||
"budget": asdict(request.budget),
|
||||
"preview_node_id": request.preview_node_id,
|
||||
"sources": {
|
||||
key: {
|
||||
"node_id": source.node_id, "batch": _batch_payload(source.batch),
|
||||
"source_ref": source.source_ref, "provider": source.provider,
|
||||
"fingerprint": source.fingerprint, "total_rows": source.total_rows,
|
||||
"truncated": source.truncated, "source_name": source.source_name,
|
||||
"kind": source.kind,
|
||||
}
|
||||
for key, source in request.sources.items()
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def _execute_reference_worker(payload: bytes) -> bytes:
|
||||
"""Data-only worker boundary; no database, provider callback or principal."""
|
||||
value = decode_worker_payload(payload, max_bytes=_TRANSPORT_BYTES)
|
||||
plan = value["plan"]
|
||||
request = BackendExecutionRequest(
|
||||
plan=ExecutionPlan(
|
||||
graph=TypedGraphIr.model_validate(plan["graph"]),
|
||||
ordered_node_ids=plan["ordered_node_ids"],
|
||||
diagnostics=tuple(DataflowDiagnostic.model_validate(item) for item in plan["diagnostics"]),
|
||||
generated_sql=plan["generated_sql"],
|
||||
sql_diagnostics=tuple(DataflowDiagnostic.model_validate(item) for item in plan["sql_diagnostics"]),
|
||||
semantic_hash=plan["semantic_hash"],
|
||||
),
|
||||
budget=ExecutionBudget(**value["budget"]),
|
||||
preview_node_id=value["preview_node_id"],
|
||||
sources={
|
||||
key: BackendSource(**{**source, "batch": _batch_from_payload(source["batch"])})
|
||||
for key, source in value["sources"].items()
|
||||
},
|
||||
)
|
||||
try:
|
||||
result = ReferenceExecutionBackend()._execute_in_process(request)
|
||||
except BackendExecutionError as exc:
|
||||
response = {"error": {
|
||||
"message": str(exc), "code": exc.code, "node_id": exc.node_id,
|
||||
"diagnostics": tuple(item.model_dump(mode="python") for item in exc.diagnostics),
|
||||
"node_diagnostics": tuple(item.model_dump(mode="python") for item in exc.node_diagnostics),
|
||||
"source_fingerprints": exc.source_fingerprints,
|
||||
"input_row_count": exc.input_row_count,
|
||||
"node_preview": exc.node_preview.model_dump(mode="python") if exc.node_preview else None,
|
||||
}}
|
||||
else:
|
||||
response = {
|
||||
"contract": result.contract.model_dump(mode="python"),
|
||||
"batch": _batch_payload(result.batch),
|
||||
"node_diagnostics": tuple(item.model_dump(mode="python") for item in result.node_diagnostics),
|
||||
"node_preview": result.node_preview.model_dump(mode="python") if result.node_preview else None,
|
||||
"metadata": dict(result.metadata),
|
||||
}
|
||||
return encode_worker_payload(response, max_bytes=_TRANSPORT_BYTES)
|
||||
|
||||
|
||||
def _validate_source_batches(request: BackendExecutionRequest) -> None:
|
||||
for source in request.sources.values():
|
||||
try:
|
||||
@@ -139,6 +297,7 @@ def _source_for_node(
|
||||
node: GraphNode,
|
||||
) -> BackendSource | None:
|
||||
candidates = (
|
||||
reference_source_key(node),
|
||||
node.id,
|
||||
str(node.config.get("source_ref") or ""),
|
||||
str(node.config.get("source_name") or ""),
|
||||
@@ -153,4 +312,10 @@ def _source_for_node(
|
||||
)
|
||||
|
||||
|
||||
def reference_source_key(node: GraphNode) -> str:
|
||||
"""Nested graphs may reuse node IDs; bind resolved data to the full config."""
|
||||
content = json.dumps(node.config, sort_keys=True, separators=(",", ":"), default=str)
|
||||
return "reference-config:" + hashlib.sha256(content.encode("utf-8")).hexdigest()
|
||||
|
||||
|
||||
__all__ = ["ReferenceExecutionBackend"]
|
||||
|
||||
@@ -19,6 +19,7 @@ from govoplan_dataflow.backend.operator_registry import (
|
||||
OperatorExecutionContext,
|
||||
OperatorExecutionResult,
|
||||
)
|
||||
from govoplan_dataflow.backend.preview_limits import MAX_RESULT_BYTES
|
||||
from govoplan_dataflow.backend.schemas import (
|
||||
DataflowDiagnostic,
|
||||
GraphNode,
|
||||
@@ -32,7 +33,6 @@ from govoplan_dataflow.backend.subflows import substitute_parameters
|
||||
|
||||
EXECUTOR_VERSION = "dataflow-preview-v2"
|
||||
MAX_EXECUTION_SECONDS = 2.0
|
||||
MAX_RESULT_BYTES = 1_000_000
|
||||
MAX_SOURCE_ROWS = 250
|
||||
MAX_INTERMEDIATE_ROWS = 10_000
|
||||
|
||||
|
||||
@@ -12,6 +12,8 @@ import sqlglot
|
||||
from sqlglot import exp
|
||||
from sqlglot.errors import ParseError
|
||||
|
||||
from govoplan_dataflow.backend.preview_limits import MAX_RESULT_BYTES
|
||||
|
||||
|
||||
ExpressionDataType = Literal[
|
||||
"unknown",
|
||||
@@ -438,6 +440,14 @@ def _evaluate_pad(expression: exp.Expression, row: dict[str, Any]) -> str | None
|
||||
target_length = int(_evaluate(expression.expression, row))
|
||||
if target_length < 0:
|
||||
raise ValueError("Padding length cannot be negative.")
|
||||
# Every character takes at least one serialized byte. Enforce the existing
|
||||
# node budget before padding allocates memory, including when an outer
|
||||
# LENGTH/SUBSTRING would otherwise conceal the oversized intermediate value.
|
||||
# The node's final byte check still accounts for Unicode and JSON overhead.
|
||||
if target_length > MAX_RESULT_BYTES:
|
||||
raise ExpressionError(
|
||||
f"Padding length exceeds the {MAX_RESULT_BYTES:,}-byte preview result limit."
|
||||
)
|
||||
source = str(value)
|
||||
if len(source) >= target_length:
|
||||
return source[:target_length]
|
||||
|
||||
@@ -1,9 +1,8 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from dataclasses import replace
|
||||
from typing import Iterable
|
||||
|
||||
from govoplan_core.core.modules import DocumentationTopic
|
||||
from govoplan_core.core.modules import DocumentationTopic, localize_documentation_topics as _localize_topics
|
||||
|
||||
|
||||
_TRANSLATIONS = {
|
||||
@@ -18,14 +17,22 @@ _TRANSLATIONS = {
|
||||
"title": "Modulgrenze von Dataflow",
|
||||
"summary": "Versionierte tabellarische Transformationen grafisch oder mit eingeschränktem SQL bearbeiten.",
|
||||
"body": (
|
||||
"Dataflow verantwortet kanonische Pipeline-Graphen, unveränderliche Revisionen, Validierung, eingeschränkte SQL-Kompilierung, Vorschau- und Laufdiagnosen sowie Herkunftsreferenzen. Datasources verantwortet den gesteuerten Katalog und Materialisierungen, Connectors den externen Abruf und Zugangsdaten, Reporting die analytische Darstellung und Exporte, Workflow die Orchestrierung und menschliche Übergaben und Risk Compliance die Sanktionsprüfung und Richtliniengrenzen. Benutzer-SQL wird in freigegebene Transformationen übersetzt und nie ungeprüft an eine Datenbank weitergereicht."
|
||||
"Dataflow verantwortet kanonische Pipeline-Graphen, unveränderliche Revisionen, Validierung, eingeschränkte SQL-Kompilierung, Vorschau- und Laufdiagnosen sowie Herkunftsreferenzen. Datasources verantwortet den gesteuerten Katalog und Materialisierungen, Connectors den externen Abruf und Zugangsdaten, Reporting die analytische Darstellung und Exporte, Workflow die Orchestrierung und menschliche Übergaben und Risk Compliance die Sanktionsprüfung und Richtliniengrenzen. Benutzer-SQL wird in freigegebene Transformationen übersetzt und nie ungeprüft an eine Datenbank weitergereicht. "
|
||||
"Das Öffnen oder Neuladen von Dataflow startet keine Pipeline. Kann der Editor nach einer Entwicklungsaktualisierung nicht geladen werden, sichern Sie ungespeicherte Arbeit vor dem Neuladen des Browsers. Administratoren sollten Fehler beim Laden von Oberflächendateien von Fehlern der Pipeline-API oder Zugriffsfehlern unterscheiden. Entwicklungs- und Browser-Konformitätsserver verwenden getrennte Abhängigkeitscaches; ein älterer Server muss nach dieser Konfigurationsaktualisierung gegebenenfalls neu gestartet werden."
|
||||
),
|
||||
},
|
||||
"dataflow.reference.nodes-and-expressions": {
|
||||
"title": "Dataflow-Knoten und Ausdrücke",
|
||||
"summary": "Typisierte Knoteneingaben, Ausdrücke, Schemafortschreibung und begrenzte Zwischenergebnisse verstehen.",
|
||||
"body": (
|
||||
"Jeder Graphknoten definiert typisierte Eingaben, Konfiguration, Ausgabeschema und Validierungsregeln. Quellknoten binden Inline-Inhalte oder gesteuerte Datasource-Referenzen; Verknüpfungs-, Filter-, Transformations-, Qualitäts-, Abgleich-, Teilfluss- und Ausgabeknoten bleiben im kanonischen Graphen ausdrücklich sichtbar. Wiederverwendbare Teilflüsse wählen eine durch Policy erlaubte unveränderliche Fluss- oder Vorlagenrevision. Der Server löst Graph, Quell-Hash, Policy-Entscheidung und geschlossene Ein-/Ausgabeverträge auf und bindet sie; mitgelieferte Graphkopien werden ignoriert, unvereinbare Eingaben und zyklische Referenzen abgelehnt. Abgleichzeilen führen stabile Schlüssel- und Eingabe-Hashes sowie Vorher-/Nachher-Werte. Prüfentscheidungen werden als unveränderliche, mandanteneigene Entscheidungssätze gespeichert; geänderte Eingaben werden ungültig, ohne Fachdaten still umzuschreiben. Ausdrücke führen weder Host- noch Datenbankcode aus. Knoten-Vorschauen sind begrenzt, für die handelnde Person datenschutzgefiltert und werden nicht als Laufergebnis gespeichert. SQL wird in denselben Graphen kompiliert; nicht unterstützte Anweisungen erscheinen als Diagnose."
|
||||
"Jeder Graphknoten definiert typisierte Eingaben, Konfiguration, Ausgabeschema und Validierungsregeln. Quellknoten binden Inline-Inhalte oder gesteuerte Datasource-Referenzen; Verknüpfungs-, Filter-, Transformations-, Qualitäts-, Abgleich-, Teilfluss- und Ausgabeknoten bleiben im kanonischen Graphen ausdrücklich sichtbar. Wiederverwendbare Teilflüsse wählen eine durch Policy erlaubte unveränderliche Fluss- oder Vorlagenrevision. Der Server löst Graph, Quell-Hash, Policy-Entscheidung und geschlossene Ein-/Ausgabeverträge auf und bindet sie; mitgelieferte Graphkopien werden ignoriert, unvereinbare Eingaben und zyklische Referenzen abgelehnt. Abgleichzeilen führen stabile Schlüssel- und Eingabe-Hashes sowie Vorher-/Nachher-Werte. Prüfentscheidungen werden als unveränderliche, mandanteneigene Entscheidungssätze gespeichert; geänderte Eingaben werden ungültig, ohne Fachdaten still umzuschreiben. Ausdrücke führen weder Host- noch Datenbankcode aus. Knoten-Vorschauen sind begrenzt, für die handelnde Person datenschutzgefiltert und werden nicht als Laufergebnis gespeichert. SQL wird in denselben Graphen kompiliert; nicht unterstützte Anweisungen erscheinen als Diagnose. "
|
||||
"Referenz-Vorschauen behalten die bestehende Grenze von 1.000.000 Byte je serialisiertem Knotenergebnis. "
|
||||
"LPAD und RPAD weisen Ziellängen über 1.000.000 Zeichen vor dem Reservieren des Auffüllspeichers zurück. "
|
||||
"Dies gilt auch für übergroße Zwischenergebnisse innerhalb von LENGTH oder SUBSTRING, selbst wenn der endgültige Einzelwert klein wäre. "
|
||||
"Verringern Sie die gewünschte Auffülllänge; der Ausdruck scheitert mit einer Diagnose am betreffenden Knoten, statt Daten abzuschneiden. "
|
||||
"Gewöhnliches Unicode-Auffüllen, NULL-Eingaben und Kürzungen innerhalb der Grenze behalten ihr Verhalten. "
|
||||
"Die abschließende Byteprüfung berücksichtigt weiterhin JSON- und Mehrbyte-Zeichenaufwand. "
|
||||
"Diese Speicherprüfung ersetzt keine Laufzeit- oder Bytegrenzen für andere Ausdrucksoperationen."
|
||||
),
|
||||
},
|
||||
"dataflow.reference.fields-and-consequences": {
|
||||
@@ -48,15 +55,4 @@ _TRANSLATIONS = {
|
||||
def localize_documentation_topics(
|
||||
topics: Iterable[DocumentationTopic],
|
||||
) -> tuple[DocumentationTopic, ...]:
|
||||
localized: list[DocumentationTopic] = []
|
||||
for topic in topics:
|
||||
german = _TRANSLATIONS.get(topic.id)
|
||||
if german is None:
|
||||
localized.append(topic)
|
||||
continue
|
||||
translations = {
|
||||
locale: dict(value) for locale, value in topic.translations.items()
|
||||
}
|
||||
translations["de"] = {**translations.get("de", {}), **german}
|
||||
localized.append(replace(topic, translations=translations))
|
||||
return tuple(localized)
|
||||
return _localize_topics(topics, locale="de", translations=_TRANSLATIONS)
|
||||
|
||||
@@ -63,7 +63,7 @@ from govoplan_dataflow.backend.german_documentation import (
|
||||
|
||||
MODULE_ID = "dataflow"
|
||||
MODULE_NAME = "Dataflow"
|
||||
MODULE_VERSION = "0.1.23"
|
||||
MODULE_VERSION = "0.1.25"
|
||||
|
||||
READ_SCOPE = "dataflow:pipeline:read"
|
||||
WRITE_SCOPE = "dataflow:pipeline:write"
|
||||
@@ -154,6 +154,118 @@ ROLE_TEMPLATES = (
|
||||
)
|
||||
|
||||
DOCUMENTATION = localize_documentation_topics((
|
||||
DocumentationTopic(
|
||||
id="dataflow.csv-source-fidelity",
|
||||
title="Import CSV without silently changing its values",
|
||||
summary="Choose text preservation or explicit legacy inference before creating a durable datasource.",
|
||||
body=(
|
||||
"The CSV import dialog defaults to Preserve text (no automatic conversion). Field whitespace, decimal digits, large identifiers, boolean-looking text and explicit empty records remain strings. "
|
||||
"Choose Infer types (legacy) only when you want the existing numeric/boolean conversion and empty-row rules. JSON imports are unchanged. Existing API clients omitting csv_value_mode retain legacy_typed behavior. "
|
||||
"The imported datasource is a deliberate durable upload, not a retained preview. Datasources owns its original UTF-8 CSV text, parsed rows, approval evidence, immutable materialization and retention; Core verifies that original input and parsed scalar types/values agree. "
|
||||
"Both original and parsed content are bounded to 5 MB and the table to 10,000 rows. Original-source export is available through the Datasources administrator API only with unrestricted current and historical visibility, never through catalogue metadata. "
|
||||
"Missing older originals cannot be reconstructed. Exact text mode can change the inferred schema to strings; review downstream numeric comparisons and conversions before adopting the new source."
|
||||
),
|
||||
layer="always", documentation_types=("user", "admin"), audience=("user", "module_admin", "operator"), order=8,
|
||||
translations={"de": {
|
||||
"title": "CSV importieren, ohne Werte unbemerkt zu ändern",
|
||||
"summary": "Vor dem dauerhaften Import zwischen Texterhalt und ausdrücklicher bisheriger Typableitung wählen.",
|
||||
"body": (
|
||||
"Der CSV-Import wählt standardmäßig Text erhalten (keine automatische Umwandlung). Leerzeichen in Werten, Dezimalstellen, große Kennungen, boolesch wirkender Text und ausdrücklich leere Datensätze bleiben Zeichenketten. "
|
||||
"Wählen Sie Typen ableiten (bisheriges Verhalten), wenn Sie die bisherige Zahlen-/Wahrheitswertumwandlung und Behandlung leerer Zeilen benötigen. JSON-Importe bleiben unverändert; API-Aufrufe ohne csv_value_mode behalten legacy_typed. "
|
||||
"Der Import ist eine bewusst dauerhafte Datenquelle und keine gespeicherte Vorschau. Datasources verantwortet Originaltext als UTF-8, verarbeitete Zeilen, Freigabenachweis, unveränderliche Materialisierung und Aufbewahrung. Core prüft die Übereinstimmung von Quelltext sowie genauen Typen und Werten. "
|
||||
"Original und verarbeiteter Inhalt sind jeweils auf 5 MB, die Tabelle auf 10.000 Zeilen begrenzt. Das Original kann nur über die Datasources-Administrations-API mit uneingeschränkter aktueller und historischer Sichtbarkeit abgerufen werden, nicht über Katalogmetadaten. "
|
||||
"Fehlende ältere Originale lassen sich nicht rekonstruieren. Der Textmodus kann Spalten zu Zeichenketten machen; prüfen Sie deshalb nachgelagerte Zahlenvergleiche und Umwandlungen vor der Übernahme."
|
||||
),
|
||||
}},
|
||||
),
|
||||
DocumentationTopic(
|
||||
id="dataflow.save-completion",
|
||||
title="Editing while a pipeline save completes",
|
||||
summary="Keep newer local edits separate from the immutable revision accepted by the server.",
|
||||
body=(
|
||||
"You may continue editing a pipeline while Save is pending. The accepted server revision becomes the saved baseline; "
|
||||
"fields changed since submission remain in the local unsaved draft and require another explicit save. Graphs are kept as whole values, "
|
||||
"not merged or reordered node by node. Save-and-leave does not navigate while newer edits remain unsaved. "
|
||||
"A second save uses the revision actually accepted by the first request; duplicate concurrent save submissions are blocked. "
|
||||
"Selecting, replacing or discarding a draft, leaving the page, or changing authentication context prevents an old completion from replacing the current editor. "
|
||||
"Such a request may already have succeeded on the server: reload and review before retrying if the context changed. "
|
||||
"Ordinary session refreshes with the same identity, credentials and permissions keep accepted IDs and revisions; cosmetic profile changes do not interrupt saving. "
|
||||
"If a save was accepted across a real authorization change, further saves in that edit session are blocked until the draft is replaced after review, so a new pipeline is not created twice. "
|
||||
"Revision conflicts retain the local draft and require review; neither the UI nor administrators automatically overwrite a conflicting server revision. "
|
||||
"Current permissions and governance remain server-enforced. No preview rows are persisted by this editor behavior."
|
||||
),
|
||||
layer="always", documentation_types=("user", "admin"), audience=("user", "module_admin", "operator"), order=7,
|
||||
translations={"de": {
|
||||
"title": "Während des Speicherns einer Pipeline weiterarbeiten",
|
||||
"summary": "Neuere lokale Änderungen von der unveränderlichen, serverseitig angenommenen Revision trennen.",
|
||||
"body": (
|
||||
"Während Speichern läuft, können Sie die Pipeline weiter bearbeiten. Die angenommene Serverrevision wird zum gespeicherten Vergleichsstand; "
|
||||
"seit dem Absenden geänderte Felder bleiben im lokalen, ungespeicherten Entwurf und benötigen einen weiteren ausdrücklichen Speichervorgang. "
|
||||
"Graphen bleiben vollständige Werte und werden nicht knotenweise zusammengeführt oder umsortiert. Speichern und Verlassen navigiert nicht, solange neuere Änderungen ungespeichert sind. "
|
||||
"Ein zweiter Speichervorgang verwendet die tatsächlich angenommene Revision des ersten; doppelte gleichzeitige Speicheranfragen werden blockiert. "
|
||||
"Auswahl, Ersetzen oder Verwerfen eines Entwurfs, Verlassen der Seite oder ein geänderter Authentifizierungskontext verhindern, dass ein altes Ergebnis den aktuellen Editor ersetzt. "
|
||||
"Die Anfrage kann auf dem Server bereits erfolgreich gewesen sein: Nach einem Kontextwechsel vor einem erneuten Versuch neu laden und prüfen. "
|
||||
"Gewöhnliche Sitzungsaktualisierungen mit gleicher Identität, gleichen Zugangsdaten und Rechten behalten angenommene IDs und Revisionen; rein optische Profiländerungen unterbrechen das Speichern nicht. "
|
||||
"Wurde ein Speichervorgang während einer tatsächlichen Berechtigungsänderung angenommen, bleiben weitere Speicheranfragen dieser Bearbeitungssitzung bis zum geprüften Ersetzen des Entwurfs gesperrt, damit keine Pipeline doppelt entsteht. "
|
||||
"Bei Revisionskonflikten bleibt der lokale Entwurf erhalten und muss geprüft werden; weder Oberfläche noch Administratoren überschreiben automatisch eine widersprechende Serverrevision. "
|
||||
"Aktuelle Rechte und Governance werden weiterhin serverseitig geprüft. Dieses Editorverhalten speichert keine Vorschauzeilen dauerhaft."
|
||||
),
|
||||
}},
|
||||
),
|
||||
DocumentationTopic(
|
||||
id="dataflow.reference-worker-limits",
|
||||
title="Reference execution process limits",
|
||||
summary="Contain expensive expressions and intermediate allocations without changing datasource authority.",
|
||||
body=(
|
||||
"Reference previews and reference development runs evaluate in a fresh disposable process, including regex and aggregate intermediate allocations. "
|
||||
"The existing row, node and per-result byte checks remain. The process additionally enforces the request's wall-clock and memory budgets "
|
||||
"(default preview: 2 seconds and 256 MiB virtual address space), rounded-up CPU seconds, no regular-file output, and 32 MiB per data-only input/result transport. "
|
||||
"Supported process budgets are at most 600 seconds and 8 GiB; unsupported controls or exceeded limits fail with structured backend.process diagnostics, never inline fallback. "
|
||||
"Datasource authorization and bounded source reads remain in the parent, including nested subflow sources; sessions and credentials are not passed to the child. "
|
||||
"Reference source collection also checks a cumulative 32 MiB typed-data budget before constructing another columnar copy; providers retain their separate per-read limits. "
|
||||
"Completed node diagnostics survive ordinary evaluation errors; a killed worker returns no partial rows or invented node progress. "
|
||||
"GOVOPLAN_ISOLATED_PROCESS_CONCURRENCY limits shared isolated-work admission per API/worker process, default 1; busy capacity is retryable. "
|
||||
"It is not a fleet-wide quota or arbitrary-code sandbox. Cancellation checks before/after reference runs remain; hard wall limits stop an unresponsive expression. "
|
||||
"Staging/production still require the separate DuckDB backend and are not converted to reference execution."
|
||||
),
|
||||
layer="static", documentation_types=("user", "admin"), audience=("user", "module_admin", "operator"), order=6,
|
||||
translations={"de": {
|
||||
"title": "Prozessgrenzen der Referenzausführung",
|
||||
"summary": "Aufwendige Ausdrücke und Zwischenspeicher begrenzen, ohne Datenquellenrechte zu verändern.",
|
||||
"body": (
|
||||
"Referenz-Vorschauen und Referenz-Entwicklungsläufe werten Ausdrücke in einem frischen, kurzlebigen Prozess aus, einschließlich regulärer Ausdrücke und großer Zwischenergebnisse. "
|
||||
"Bestehende Zeilen-, Knoten- und Ergebnis-Bytegrenzen bleiben bestehen. Zusätzlich gelten das Laufzeit- und Speicherbudget der Anfrage "
|
||||
"(Vorschau standardmäßig 2 Sekunden und 256 MiB virtueller Adressraum), aufgerundete CPU-Sekunden, keine regulären Ausgabedateien und je 32 MiB für den reinen Datentransport. "
|
||||
"Prozessbudgets unterstützen höchstens 600 Sekunden und 8 GiB. Fehlende Betriebssystemkontrollen oder überschrittene Grenzen erzeugen strukturierte backend.process-Diagnosen, ohne Ausweichbetrieb im Hauptprozess. "
|
||||
"Datenquellenrechte und begrenzte Quellabrufe werden im Hauptprozess geprüft, auch für verschachtelte Teilflüsse; Sitzungen und Zugangsdaten gelangen nicht in den Kindprozess. "
|
||||
"Referenz-Quellabrufe prüfen außerdem zusammen höchstens 32 MiB typisierte Daten, bevor eine weitere spaltenweise Kopie entsteht; getrennte Abrufgrenzen der Anbieter bleiben bestehen. "
|
||||
"Gewöhnliche Auswertungsfehler behalten bereits abgeschlossene Knotendiagnosen. Ein gestoppter Prozess liefert keine Teilzeilen und keinen erfundenen Knotenfortschritt. "
|
||||
"GOVOPLAN_ISOLATED_PROCESS_CONCURRENCY begrenzt gemeinsam genutzte isolierte Arbeit je API-/Worker-Prozess, standardmäßig 1; bei Auslastung ist ein erneuter Versuch möglich. "
|
||||
"Dies ist weder eine systemweite Quote noch eine Sandbox für beliebigen Code. Abbruchprüfungen vor und nach Referenzläufen bleiben erhalten; harte Laufzeitgrenzen stoppen hängende Ausdrücke. "
|
||||
"Staging und Produktion benötigen weiterhin das gesonderte DuckDB-Backend und wechseln nicht zur Referenzausführung."
|
||||
),
|
||||
}},
|
||||
),
|
||||
DocumentationTopic(
|
||||
id="dataflow.workspace-layout",
|
||||
title="Dataflow workspace actions",
|
||||
summary="Find collection-wide commands in their consistent workspace position.",
|
||||
body="The workspace documentation book sits beside Pipelines; automation and run help sits beside "
|
||||
"the corresponding dialog title, and field help stays with its label. "
|
||||
"Reload and New pipeline use the persistent full-width workspace header at the upper right; Reload sits immediately before creation. Selecting a record, changing filters, or opening an editor does not move these collection-wide commands into the left pane. SQL editing, validation, previews, saving, and execution keep their existing editor scope and bounded safety rules. Existing permissions, disabled-state rules, and unsaved-change guards still apply. Administrators configure authority through the existing permission system; no new permission or automatic operation is introduced.",
|
||||
layer="static",
|
||||
documentation_types=("user", "admin"),
|
||||
audience=("user", "module_admin", "operator"),
|
||||
order=5,
|
||||
translations={"de": {
|
||||
"title": "Datenflüsse: Aktionen im Arbeitsbereich",
|
||||
"summary": "Sammlungsweite Aktionen an ihrer einheitlichen Position im Arbeitsbereich finden.",
|
||||
"body": "Das Dokumentationsbuch des Arbeitsbereichs steht neben Pipelines; Hilfe zu Automatisierung "
|
||||
"und Ausführung steht neben dem jeweiligen Dialogtitel, und Feldhilfe bleibt bei der "
|
||||
"Feldbezeichnung. "
|
||||
"Neu laden und Neue Pipeline stehen oben rechts in der dauerhaft sichtbaren, arbeitsbereichsweiten Leiste; Neu laden steht unmittelbar vor dem Anlegen. Auswahl, Filterwechsel und Bearbeitung verschieben diese sammlungsweiten Aktionen nicht in den linken Bereich. SQL-Bearbeitung, Validierung, Vorschau, Speichern und Ausführung behalten ihren bisherigen Editorbereich und ihre begrenzenden Sicherheitsregeln. Bestehende Berechtigungen, Deaktivierungsregeln und der Schutz ungespeicherter Änderungen gelten weiterhin. Administratoren konfigurieren Rechte im bestehenden Berechtigungssystem; es entstehen weder neue Rechte noch automatische Vorgänge.",
|
||||
}},
|
||||
),
|
||||
DocumentationTopic(
|
||||
id="dataflow.data-subject-requests",
|
||||
title="Dataflow data-subject requests",
|
||||
@@ -186,7 +298,12 @@ DOCUMENTATION = localize_documentation_topics((
|
||||
"and credentials; Reporting owns analytical presentation and exports; "
|
||||
"Workflow owns orchestration and human handoffs; Risk Compliance owns sanctions review "
|
||||
"semantics and policy gates. User SQL is compiled into approved transforms and is never "
|
||||
"passed unchecked to a backing database."
|
||||
"passed unchecked to a backing database. Opening or reloading Dataflow does not "
|
||||
"start a pipeline. If the editor cannot be loaded after a development update, "
|
||||
"preserve unsaved work before reloading the browser. Administrators should "
|
||||
"distinguish frontend asset failures from pipeline API or access errors. "
|
||||
"Development and browser-conformance servers use separate dependency caches; "
|
||||
"an older server may need restarting after this configuration update."
|
||||
),
|
||||
layer="available",
|
||||
documentation_types=("admin", "user"),
|
||||
@@ -280,7 +397,14 @@ DOCUMENTATION = localize_documentation_topics((
|
||||
"Expressions use the typed Dataflow expression language and never execute arbitrary host or database "
|
||||
"code. Selecting a node may request a bounded intermediate preview; preview rows are transient, "
|
||||
"privacy-filtered for the actor, and are not retained as run output. SQL editing compiles into the same "
|
||||
"canonical graph, so unsupported statements are diagnostics rather than pass-through SQL."
|
||||
"canonical graph, so unsupported statements are diagnostics rather than pass-through SQL. "
|
||||
"Reference previews retain the existing 1,000,000-byte serialized result limit per node. "
|
||||
"LPAD and RPAD reject target lengths above 1,000,000 characters before allocating padding, "
|
||||
"including oversized intermediate values inside LENGTH or SUBSTRING even if the final scalar would be small. "
|
||||
"Reduce the requested padding length; this fails the expression at its node instead of truncating data. "
|
||||
"Ordinary Unicode padding, null inputs, and in-budget truncation keep their existing behavior; "
|
||||
"the final byte check still accounts for JSON and multibyte character overhead. "
|
||||
"This allocation guard does not replace runtime or byte limits for other expression operations."
|
||||
),
|
||||
layer="available",
|
||||
documentation_types=("admin", "user"),
|
||||
|
||||
@@ -0,0 +1,3 @@
|
||||
"""Shared existing limits for reference-preview results and allocating expressions."""
|
||||
|
||||
MAX_RESULT_BYTES = 1_000_000
|
||||
@@ -32,6 +32,8 @@ from govoplan_core.core.references import (
|
||||
validate_access_scope_reference,
|
||||
)
|
||||
from govoplan_core.core.tabular_sources import (
|
||||
TabularCsvSource,
|
||||
TabularSourceError,
|
||||
parse_tabular_csv,
|
||||
)
|
||||
from govoplan_core.db.session import get_session
|
||||
@@ -432,6 +434,7 @@ def api_create_source_snapshot(
|
||||
payload.csv_text or "",
|
||||
delimiter=payload.delimiter,
|
||||
max_rows=10_000,
|
||||
value_mode=payload.csv_value_mode,
|
||||
)
|
||||
if payload.format == "csv"
|
||||
else tuple(payload.rows or ())
|
||||
@@ -448,6 +451,11 @@ def api_create_source_snapshot(
|
||||
shape="tabular",
|
||||
rows=rows,
|
||||
provider="dataflow.upload",
|
||||
csv_source=(TabularCsvSource(
|
||||
text=payload.csv_text or "",
|
||||
delimiter=payload.delimiter,
|
||||
value_mode=payload.csv_value_mode,
|
||||
) if payload.format == "csv" else None),
|
||||
provenance={
|
||||
"created_via": "dataflow",
|
||||
"source_format": payload.format,
|
||||
@@ -462,6 +470,8 @@ def api_create_source_snapshot(
|
||||
principal,
|
||||
stage_ref=stage.ref,
|
||||
)
|
||||
except TabularSourceError as exc:
|
||||
raise HTTPException(status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc)) from exc
|
||||
except DatasourceError as exc:
|
||||
raise _source_http_error(exc) from exc
|
||||
audit_event(
|
||||
|
||||
@@ -737,6 +737,7 @@ class TabularSnapshotCreateRequest(BaseModel):
|
||||
format: Literal["json", "csv"] = "json"
|
||||
rows: list[dict[str, Any]] | None = Field(default=None, max_length=10_000)
|
||||
csv_text: str | None = Field(default=None, max_length=5_000_000)
|
||||
csv_value_mode: Literal["legacy_typed", "text"] = "legacy_typed"
|
||||
delimiter: str = Field(default=",", min_length=1, max_length=1)
|
||||
|
||||
@model_validator(mode="after")
|
||||
|
||||
@@ -28,6 +28,7 @@ from govoplan_core.core.datasources import (
|
||||
datasource_publication,
|
||||
)
|
||||
from govoplan_core.db.base import utcnow
|
||||
from govoplan_core.security.worker_payload import WorkerPayloadError, encode_worker_payload
|
||||
from govoplan_dataflow.backend.backends import (
|
||||
BackendExecutionError,
|
||||
BackendSource,
|
||||
@@ -35,6 +36,7 @@ from govoplan_dataflow.backend.backends import (
|
||||
execute_typed_graph,
|
||||
)
|
||||
from govoplan_dataflow.backend.batches import TypedBatch
|
||||
from govoplan_dataflow.backend.backends.reference import reference_source_key
|
||||
from govoplan_dataflow.backend.db.models import (
|
||||
DataflowPipeline,
|
||||
DataflowPipelineDeployment,
|
||||
@@ -48,7 +50,6 @@ from govoplan_dataflow.backend.executor import (
|
||||
PipelineExecutionError,
|
||||
PipelineExecutionResult,
|
||||
ResolvedSource,
|
||||
execute_preview,
|
||||
)
|
||||
from govoplan_dataflow.backend.governance import (
|
||||
definition_governance_payload,
|
||||
@@ -1122,20 +1123,11 @@ def _execute_pipeline_preview(
|
||||
principal=principal,
|
||||
registry=registry,
|
||||
)
|
||||
if backend == "reference":
|
||||
return (
|
||||
execute_preview(
|
||||
graph,
|
||||
row_limit=row_limit,
|
||||
source_resolver=source_resolver,
|
||||
preview_node_id=preview_node_id,
|
||||
),
|
||||
EXECUTOR_VERSION,
|
||||
)
|
||||
sources = _typed_backend_sources(
|
||||
graph,
|
||||
source_resolver=source_resolver,
|
||||
source_limit=max(MAX_SOURCE_ROWS, row_limit),
|
||||
source_limit=MAX_SOURCE_ROWS if backend == "reference" else max(MAX_SOURCE_ROWS, row_limit),
|
||||
include_subflows=backend == "reference",
|
||||
)
|
||||
try:
|
||||
result = execute_typed_graph(
|
||||
@@ -1149,8 +1141,14 @@ def _execute_pipeline_preview(
|
||||
raise PipelineExecutionError(
|
||||
str(exc),
|
||||
node_id=exc.node_id,
|
||||
diagnostics=tuple(exc.diagnostics),
|
||||
retryable=exc.code == "backend.capacity",
|
||||
diagnostics=(*exc.diagnostics, DataflowDiagnostic(
|
||||
severity="error", code=exc.code, message=str(exc), node_id=exc.node_id,
|
||||
)),
|
||||
node_diagnostics=exc.node_diagnostics,
|
||||
source_fingerprints=exc.source_fingerprints,
|
||||
input_row_count=exc.input_row_count,
|
||||
node_preview=exc.node_preview,
|
||||
retryable=exc.code in {"backend.capacity", "backend.process.busy"},
|
||||
) from exc
|
||||
columns = [
|
||||
PreviewColumn(
|
||||
@@ -1227,13 +1225,51 @@ def _typed_backend_sources(
|
||||
*,
|
||||
source_resolver,
|
||||
source_limit: int = MAX_SOURCE_ROWS,
|
||||
include_subflows: bool = False,
|
||||
_depth: int = 0,
|
||||
_remaining_source_bytes: list[int] | None = None,
|
||||
) -> dict[str, BackendSource]:
|
||||
if _depth > 5:
|
||||
raise PipelineExecutionError("Subflows are limited to five nested levels.")
|
||||
if _remaining_source_bytes is None:
|
||||
_remaining_source_bytes = [32 * 1024 * 1024]
|
||||
sources: dict[str, BackendSource] = {}
|
||||
for node in graph.nodes:
|
||||
if include_subflows and node.type == "subflow":
|
||||
from govoplan_dataflow.backend.subflows import substitute_parameters
|
||||
|
||||
parameters = node.config.get("parameters")
|
||||
nested = PipelineGraph.model_validate(substitute_parameters(
|
||||
node.config.get("graph"), parameters if isinstance(parameters, dict) else {},
|
||||
))
|
||||
sources.update(_typed_backend_sources(
|
||||
nested, source_resolver=source_resolver, source_limit=source_limit,
|
||||
include_subflows=True, _depth=_depth + 1,
|
||||
_remaining_source_bytes=_remaining_source_bytes,
|
||||
))
|
||||
if node.type != "source.reference":
|
||||
continue
|
||||
if include_subflows and _remaining_source_bytes[0] <= 0:
|
||||
raise PipelineExecutionError(
|
||||
"Combined source data exceeds the 32 MiB transfer budget.", node_id=node.id,
|
||||
)
|
||||
resolved = source_resolver(node, source_limit)
|
||||
sources[node.id] = BackendSource(
|
||||
if include_subflows:
|
||||
try:
|
||||
# Check before constructing another columnar copy. The provider
|
||||
# still owns bounds on its individual authorized read; do not keep
|
||||
# accumulating individually valid batches before the worker gate.
|
||||
encoded_size = len(encode_worker_payload(
|
||||
tuple(dict(row) for row in resolved.rows), max_bytes=_remaining_source_bytes[0],
|
||||
))
|
||||
except WorkerPayloadError as exc:
|
||||
raise PipelineExecutionError(
|
||||
"Combined source data exceeds the 32 MiB transfer budget or contains unsupported values.",
|
||||
node_id=node.id,
|
||||
) from exc
|
||||
_remaining_source_bytes[0] -= encoded_size
|
||||
key = reference_source_key(node) if include_subflows else node.id
|
||||
sources[key] = BackendSource(
|
||||
node_id=node.id,
|
||||
batch=TypedBatch.from_rows(resolved.rows),
|
||||
source_ref=resolved.source_ref,
|
||||
|
||||
Executable
+121
@@ -0,0 +1,121 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from types import SimpleNamespace
|
||||
import unittest
|
||||
from unittest.mock import Mock, patch
|
||||
|
||||
from fastapi import HTTPException
|
||||
from pydantic import ValidationError
|
||||
|
||||
from govoplan_core.auth import ApiPrincipal
|
||||
from govoplan_core.core.access import PrincipalRef
|
||||
from govoplan_dataflow.backend.router import WRITE_SCOPE, api_create_source_snapshot
|
||||
from govoplan_dataflow.backend.schemas import TabularSnapshotCreateRequest
|
||||
|
||||
|
||||
class CsvStagingRouteTests(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.principal = ApiPrincipal(
|
||||
principal=PrincipalRef(
|
||||
account_id="account",
|
||||
membership_id="membership",
|
||||
tenant_id="tenant",
|
||||
scopes=frozenset({WRITE_SCOPE}),
|
||||
),
|
||||
user=object(),
|
||||
account=object(),
|
||||
)
|
||||
self.session = Mock()
|
||||
self.writer = Mock()
|
||||
self.writer.create_stage.return_value = SimpleNamespace(ref="stage:one")
|
||||
self.writer.promote_stage.return_value = (
|
||||
SimpleNamespace(
|
||||
ref="datasource:one",
|
||||
provider="dataflow.upload",
|
||||
source_name="upload",
|
||||
row_count=1,
|
||||
fingerprint="f" * 64,
|
||||
),
|
||||
SimpleNamespace(ref="materialization:one"),
|
||||
)
|
||||
for name, options in (
|
||||
("get_registry", {"return_value": object()}),
|
||||
("datasource_lifecycle", {"return_value": self.writer}),
|
||||
("audit_event", {}),
|
||||
("_source_response", {"side_effect": lambda value: value}),
|
||||
):
|
||||
active = patch(f"govoplan_dataflow.backend.router.{name}", **options)
|
||||
active.start()
|
||||
self.addCleanup(active.stop)
|
||||
|
||||
def test_csv_stage_receives_exact_original_and_selected_parser_mode(self) -> None:
|
||||
text = 'value\r\n" 001 "\r\n'
|
||||
for mode, expected in (("text", " 001 "), ("legacy_typed", "001")):
|
||||
with self.subTest(mode=mode):
|
||||
api_create_source_snapshot(
|
||||
TabularSnapshotCreateRequest(
|
||||
name="Upload",
|
||||
source_name="upload",
|
||||
format="csv",
|
||||
csv_text=text,
|
||||
csv_value_mode=mode,
|
||||
),
|
||||
session=self.session,
|
||||
principal=self.principal,
|
||||
)
|
||||
stage = self.writer.create_stage.call_args.kwargs["stage"]
|
||||
self.assertEqual(({"value": expected},), stage.rows)
|
||||
self.assertEqual(text, stage.csv_source.text)
|
||||
self.assertEqual(mode, stage.csv_source.value_mode)
|
||||
self.assertEqual("core.csv.v1", stage.csv_source.parser_profile)
|
||||
self.assertNotIn("text", stage.metadata)
|
||||
self.assertEqual(2, self.writer.promote_stage.call_count)
|
||||
self.assertEqual(2, self.session.commit.call_count)
|
||||
|
||||
def test_json_stage_does_not_invent_csv_evidence(self) -> None:
|
||||
api_create_source_snapshot(
|
||||
TabularSnapshotCreateRequest(
|
||||
name="Upload", source_name="upload", rows=[{"value": "001"}]
|
||||
),
|
||||
session=self.session,
|
||||
principal=self.principal,
|
||||
)
|
||||
stage = self.writer.create_stage.call_args.kwargs["stage"]
|
||||
self.assertIsNone(stage.csv_source)
|
||||
self.assertEqual(({"value": "001"},), stage.rows)
|
||||
|
||||
def test_malformed_text_csv_is_422_before_any_durable_write(self) -> None:
|
||||
with self.assertRaises(HTTPException) as raised:
|
||||
api_create_source_snapshot(
|
||||
TabularSnapshotCreateRequest(
|
||||
name="Upload",
|
||||
source_name="upload",
|
||||
format="csv",
|
||||
csv_text="a,b\nonly-one\n",
|
||||
csv_value_mode="text",
|
||||
),
|
||||
session=self.session,
|
||||
principal=self.principal,
|
||||
)
|
||||
self.assertEqual(422, raised.exception.status_code)
|
||||
self.writer.create_stage.assert_not_called()
|
||||
self.writer.promote_stage.assert_not_called()
|
||||
self.session.commit.assert_not_called()
|
||||
|
||||
def test_invalid_unicode_is_rejected_by_request_schema_before_any_durable_write(
|
||||
self,
|
||||
) -> None:
|
||||
with self.assertRaises(ValidationError):
|
||||
api_create_source_snapshot(
|
||||
TabularSnapshotCreateRequest(
|
||||
name="Upload",
|
||||
source_name="upload",
|
||||
format="csv",
|
||||
csv_text="value\nprivate-\ud800\n",
|
||||
),
|
||||
session=self.session,
|
||||
principal=self.principal,
|
||||
)
|
||||
self.writer.create_stage.assert_not_called()
|
||||
self.writer.promote_stage.assert_not_called()
|
||||
self.session.commit.assert_not_called()
|
||||
@@ -67,12 +67,13 @@ class DataflowInterfaceDocumentationContractTests(unittest.TestCase):
|
||||
|
||||
for component in (
|
||||
"ActionBlockerHint",
|
||||
"DocumentationHelpLink",
|
||||
"titleHelp={<DocumentationHelpLink reference={DATAFLOW_DOCUMENTATION} />}",
|
||||
"useUnsavedDraftGuard",
|
||||
"ConfirmDialog",
|
||||
):
|
||||
self.assertIn(component, page)
|
||||
self.assertIn("DATAFLOW_NODE_DOCUMENTATION", inspector)
|
||||
self.assertNotIn("helpAction=", page)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
@@ -0,0 +1,165 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import unittest
|
||||
|
||||
from govoplan_dataflow.backend.executor import (
|
||||
MAX_RESULT_BYTES,
|
||||
PipelineExecutionError,
|
||||
execute_preview,
|
||||
)
|
||||
from govoplan_dataflow.backend.expressions import ExpressionError, evaluate_expression
|
||||
from govoplan_dataflow.backend.manifest import get_manifest
|
||||
from govoplan_dataflow.backend.schemas import (
|
||||
GraphEdge,
|
||||
GraphNode,
|
||||
GraphPosition,
|
||||
PipelineGraph,
|
||||
)
|
||||
|
||||
|
||||
class _UnreadableFill:
|
||||
def __str__(self) -> str:
|
||||
raise AssertionError(
|
||||
"Oversized padding must be rejected before reading or multiplying its fill."
|
||||
)
|
||||
|
||||
|
||||
def _expression_graph(expression: str, *, result_type: str = "string") -> PipelineGraph:
|
||||
return PipelineGraph(
|
||||
nodes=[
|
||||
GraphNode(
|
||||
id="source",
|
||||
type="source.inline",
|
||||
label="Source",
|
||||
position=GraphPosition(x=0, y=0),
|
||||
config={"source_name": "fixture", "rows": [{"value": "x"}]},
|
||||
),
|
||||
GraphNode(
|
||||
id="padding",
|
||||
type="expression",
|
||||
label="Padding",
|
||||
position=GraphPosition(x=200, y=0),
|
||||
config={
|
||||
"target_column": "padded",
|
||||
"expression": expression,
|
||||
"result_type": result_type,
|
||||
},
|
||||
),
|
||||
GraphNode(
|
||||
id="output",
|
||||
type="output",
|
||||
label="Output",
|
||||
position=GraphPosition(x=400, y=0),
|
||||
config={},
|
||||
),
|
||||
],
|
||||
edges=[
|
||||
GraphEdge(id="source-padding", source="source", target="padding"),
|
||||
GraphEdge(id="padding-output", source="padding", target="output"),
|
||||
],
|
||||
)
|
||||
|
||||
|
||||
class PaddingBudgetTests(unittest.TestCase):
|
||||
def test_oversized_padding_is_rejected_before_fill_evaluation_or_allocation(
|
||||
self,
|
||||
) -> None:
|
||||
for operation in ("lpad", "rpad"):
|
||||
for length in (MAX_RESULT_BYTES + 1, 10**100):
|
||||
with self.subTest(operation=operation, length=length):
|
||||
with self.assertRaisesRegex(
|
||||
ExpressionError, "Padding length.*1,000,000"
|
||||
):
|
||||
evaluate_expression(
|
||||
f"{operation}('x', {length}, fill)",
|
||||
{"fill": _UnreadableFill()},
|
||||
)
|
||||
|
||||
def test_budget_cannot_be_bypassed_by_hiding_large_padding_in_a_small_scalar(
|
||||
self,
|
||||
) -> None:
|
||||
for wrapper in ("length({})", "substring({}, 1, 1)"):
|
||||
with self.subTest(wrapper=wrapper):
|
||||
expression = wrapper.format(f"lpad('x', {MAX_RESULT_BYTES + 1}, fill)")
|
||||
with self.assertRaises(ExpressionError):
|
||||
evaluate_expression(expression, {"fill": _UnreadableFill()})
|
||||
|
||||
def test_existing_boundary_and_ordinary_padding_are_preserved(self) -> None:
|
||||
self.assertEqual(1_000_000, MAX_RESULT_BYTES)
|
||||
for operation in ("lpad", "rpad"):
|
||||
with self.subTest(operation=operation):
|
||||
value = evaluate_expression(
|
||||
f"{operation}('x', {MAX_RESULT_BYTES}, '0')", {}
|
||||
)
|
||||
self.assertEqual(MAX_RESULT_BYTES, len(value))
|
||||
self.assertEqual(1, value.count("x"))
|
||||
self.assertEqual(
|
||||
"abc", evaluate_expression(f"{operation}('abcdef', 3, '')", {})
|
||||
)
|
||||
self.assertEqual(
|
||||
"", evaluate_expression(f"{operation}('abcdef', 0, '')", {})
|
||||
)
|
||||
self.assertEqual(
|
||||
"abc", evaluate_expression(f"{operation}('abc', 3, '')", {})
|
||||
)
|
||||
|
||||
def test_null_negative_and_empty_fill_semantics_are_unchanged(self) -> None:
|
||||
for operation in ("lpad", "rpad"):
|
||||
with self.subTest(operation=operation):
|
||||
self.assertIsNone(
|
||||
evaluate_expression(f"{operation}(NULL, {10**100}, '')", {})
|
||||
)
|
||||
self.assertIsNone(evaluate_expression(f"{operation}(NULL, -1, '')", {}))
|
||||
with self.assertRaisesRegex(ValueError, "cannot be negative"):
|
||||
evaluate_expression(f"{operation}('x', -1, '0')", {})
|
||||
with self.assertRaisesRegex(ValueError, "fill text cannot be empty"):
|
||||
evaluate_expression(f"{operation}('x', 2, '')", {})
|
||||
|
||||
def test_multibyte_fill_and_truncation_preserve_character_semantics(self) -> None:
|
||||
self.assertEqual("ö🙂öÄ", evaluate_expression("lpad('Ä', 4, 'ö🙂')", {}))
|
||||
self.assertEqual("Äö🙂ö", evaluate_expression("rpad('Ä', 4, 'ö🙂')", {}))
|
||||
self.assertEqual("🙂ä", evaluate_expression("lpad('🙂ä中', 2, '0')", {}))
|
||||
result = execute_preview(
|
||||
_expression_graph("rpad(value, 4, 'ö🙂')"), row_limit=10
|
||||
)
|
||||
self.assertEqual("xö🙂ö", result.rows[0]["padded"])
|
||||
|
||||
def test_preview_reports_padding_guard_at_owning_node_and_retains_final_byte_limit(
|
||||
self,
|
||||
) -> None:
|
||||
with self.assertRaisesRegex(PipelineExecutionError, "Padding length") as raised:
|
||||
execute_preview(
|
||||
_expression_graph(
|
||||
f"length(lpad(value, {MAX_RESULT_BYTES + 1}, '0'))",
|
||||
result_type="integer",
|
||||
),
|
||||
row_limit=10,
|
||||
)
|
||||
self.assertEqual("padding", raised.exception.node_id)
|
||||
# Non-ASCII characters need several serialized bytes. The preallocation
|
||||
# character bound supplements, and never replaces, the node byte bound.
|
||||
with self.assertRaisesRegex(
|
||||
PipelineExecutionError, "one-megabyte result limit"
|
||||
) as raised:
|
||||
execute_preview(_expression_graph("rpad(value, 200000, 'ö')"), row_limit=10)
|
||||
self.assertEqual("padding", raised.exception.node_id)
|
||||
|
||||
def test_user_and_operator_documentation_explains_intermediate_padding_limit_in_both_languages(
|
||||
self,
|
||||
) -> None:
|
||||
topic = next(
|
||||
topic
|
||||
for topic in get_manifest().documentation
|
||||
if topic.id == "dataflow.reference.nodes-and-expressions"
|
||||
)
|
||||
self.assertIn("user", topic.documentation_types)
|
||||
self.assertIn("admin", topic.documentation_types)
|
||||
for text in (topic.body, topic.translations["de"]["body"]):
|
||||
self.assertIn("LPAD", text)
|
||||
self.assertIn("RPAD", text)
|
||||
self.assertIn("LENGTH", text)
|
||||
self.assertIn("SUBSTRING", text)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
@@ -0,0 +1,117 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import date, datetime, timezone
|
||||
from decimal import Decimal
|
||||
import time
|
||||
import unittest
|
||||
from unittest.mock import patch
|
||||
|
||||
from govoplan_core.security.bounded_process import ProcessBudgetError
|
||||
from govoplan_dataflow.backend.backends import (
|
||||
BackendExecutionError, BackendSource, ExecutionBudget, execute_typed_graph,
|
||||
)
|
||||
from govoplan_dataflow.backend.backends.reference import ReferenceExecutionBackend
|
||||
from govoplan_dataflow.backend.batches import TypedBatch
|
||||
from govoplan_dataflow.backend.executor import PipelineExecutionError, ResolvedSource
|
||||
from govoplan_dataflow.backend.manifest import get_manifest
|
||||
from govoplan_dataflow.backend.schemas import GraphEdge, GraphNode, GraphPosition, PipelineGraph
|
||||
from govoplan_dataflow.backend.service import _execute_pipeline_preview, _typed_backend_sources
|
||||
|
||||
|
||||
def graph_for(expression: str = "value", rows: list | None = None) -> PipelineGraph:
|
||||
nodes = [
|
||||
GraphNode(id="source", type="source.inline", label="Source", position=GraphPosition(x=0, y=0),
|
||||
config={"source_name": "records", "rows": rows or [{"value": "normal"}]}),
|
||||
GraphNode(id="expression", type="expression", label="Expression", position=GraphPosition(x=100, y=0),
|
||||
config={"target_column": "result", "expression": expression, "result_type": "unknown"}),
|
||||
GraphNode(id="output", type="output", label="Output", position=GraphPosition(x=200, y=0), config={}),
|
||||
]
|
||||
return PipelineGraph(nodes=nodes, edges=[
|
||||
GraphEdge(id="first", source="source", target="expression"),
|
||||
GraphEdge(id="second", source="expression", target="output"),
|
||||
])
|
||||
|
||||
|
||||
class ReferenceProcessTests(unittest.TestCase):
|
||||
def test_non_finite_deadlines_are_rejected_before_registry_wait(self) -> None:
|
||||
for value in (float("nan"), float("inf")):
|
||||
with self.subTest(value=value), self.assertRaisesRegex(ValueError, "finite"):
|
||||
ExecutionBudget(max_wall_seconds=value)
|
||||
|
||||
def test_real_child_not_parent_helper_and_typed_result_preserved(self) -> None:
|
||||
row = {"value": Decimal("1.20"), "date": date(2026, 9, 8),
|
||||
"when": datetime(2026, 9, 8, tzinfo=timezone.utc), "binary": b"\x00\xff"}
|
||||
graph = graph_for()
|
||||
graph.nodes[0] = graph.nodes[0].model_copy(update={"type": "source.reference", "config": {
|
||||
"source_ref": "datasource:fixture", "source_name": "records",
|
||||
}})
|
||||
source = BackendSource(node_id="source", batch=TypedBatch.from_rows([row]),
|
||||
source_ref="datasource:fixture", provider="test", fingerprint="pinned", total_rows=1)
|
||||
with patch.object(ReferenceExecutionBackend, "_execute_in_process", side_effect=AssertionError("parent evaluation")):
|
||||
result = execute_typed_graph(graph, backend="reference", sources={"source": source})
|
||||
self.assertEqual(result.rows, [{**row, "result": Decimal("1.20")}])
|
||||
self.assertEqual(result.contract.lineage.source_fingerprints[0]["fingerprint"], "pinned")
|
||||
|
||||
def test_real_pathological_regex_is_stopped_with_structured_failure(self) -> None:
|
||||
graph = graph_for("regexp_full_match(value, '(a+)+$')", [{"value": "a" * 100 + "!"}])
|
||||
started = time.monotonic()
|
||||
with self.assertRaises(BackendExecutionError) as caught:
|
||||
execute_typed_graph(graph, backend="reference", budget=ExecutionBudget(max_wall_seconds=2))
|
||||
self.assertIn(caught.exception.code, {"backend.process.timeout", "backend.process.cpu_limit"})
|
||||
self.assertLess(time.monotonic() - started, 4)
|
||||
|
||||
def test_aggregate_padding_allocation_is_contained_by_child_memory_limit(self) -> None:
|
||||
graph = graph_for("lpad(value, 900000, '0')", [{"value": "x"} for _ in range(250)])
|
||||
with self.assertRaises(BackendExecutionError) as caught:
|
||||
execute_typed_graph(graph, backend="reference", budget=ExecutionBudget(
|
||||
max_wall_seconds=5, max_memory_bytes=128 * 1024 * 1024,
|
||||
))
|
||||
self.assertEqual(caught.exception.code, "backend.process.memory_limit")
|
||||
|
||||
def test_busy_preview_is_retryable_and_never_evaluates_inline(self) -> None:
|
||||
with patch("govoplan_dataflow.backend.backends.reference.run_bounded_operation", side_effect=ProcessBudgetError("busy")):
|
||||
with self.assertRaises(PipelineExecutionError) as caught:
|
||||
_execute_pipeline_preview(graph_for(), session=None, principal=None, registry=None,
|
||||
backend="reference", row_limit=10, preview_node_id=None)
|
||||
self.assertTrue(caught.exception.retryable)
|
||||
self.assertEqual(caught.exception.diagnostics[-1].code, "backend.process.busy")
|
||||
|
||||
def test_nested_source_ids_do_not_alias_different_authorized_data(self) -> None:
|
||||
outer = graph_for()
|
||||
inner = graph_for()
|
||||
for graph, ref in ((outer, "datasource:outer"), (inner, "datasource:inner")):
|
||||
graph.nodes[0] = graph.nodes[0].model_copy(update={"type": "source.reference", "config": {
|
||||
"source_ref": ref, "source_name": "records",
|
||||
}})
|
||||
outer.nodes[1] = outer.nodes[1].model_copy(update={"type": "subflow", "config": {
|
||||
"graph": inner.model_dump(mode="python"), "parameters": {},
|
||||
}})
|
||||
def resolve(node, limit):
|
||||
return ResolvedSource(rows=({"value": node.config["source_ref"]},),
|
||||
source_ref=node.config["source_ref"], provider="test",
|
||||
fingerprint=node.config["source_ref"], total_rows=1)
|
||||
sources = _typed_backend_sources(outer, source_resolver=resolve, include_subflows=True)
|
||||
self.assertEqual({source.source_ref for source in sources.values()}, {"datasource:outer", "datasource:inner"})
|
||||
self.assertEqual(len(sources), 2)
|
||||
|
||||
def test_static_worker_documentation_is_bilingual(self) -> None:
|
||||
topic = next(item for item in get_manifest().documentation if item.id == "dataflow.reference-worker-limits")
|
||||
for body in (topic.body, topic.translations["de"]["body"]):
|
||||
self.assertIn("GOVOPLAN_ISOLATED_PROCESS_CONCURRENCY", body)
|
||||
self.assertIn("32 MiB", body)
|
||||
self.assertEqual(set(topic.documentation_types), {"user", "admin"})
|
||||
|
||||
def test_cumulative_source_budget_stops_before_reading_further_sources(self) -> None:
|
||||
graph = graph_for()
|
||||
graph.nodes = [graph.nodes[0].model_copy(update={
|
||||
"id": f"source-{index}", "type": "source.reference", "config": {"source_ref": f"fixture:{index}"},
|
||||
}) for index in range(3)]
|
||||
calls = []
|
||||
def resolve(node, limit):
|
||||
calls.append(node.id)
|
||||
return ResolvedSource(rows=({"value": "x" * 100},), source_ref=node.config["source_ref"],
|
||||
provider="test", fingerprint="fixed", total_rows=1)
|
||||
with self.assertRaisesRegex(PipelineExecutionError, "Combined source data"):
|
||||
_typed_backend_sources(graph, source_resolver=resolve, include_subflows=True,
|
||||
_remaining_source_bytes=[200])
|
||||
self.assertEqual(calls, ["source-0", "source-1"])
|
||||
@@ -0,0 +1,194 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
from types import SimpleNamespace
|
||||
import unittest
|
||||
from unittest.mock import Mock, patch
|
||||
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from govoplan_core.auth import ApiPrincipal
|
||||
from govoplan_core.core.access import PrincipalRef
|
||||
from govoplan_core.core.datasources import (
|
||||
CAPABILITY_DATASOURCE_CATALOGUE,
|
||||
DatasourceAccessError,
|
||||
DatasourceCatalogueProvider,
|
||||
DatasourceDescriptor,
|
||||
DatasourceReadRequest,
|
||||
DatasourceReadResult,
|
||||
)
|
||||
from govoplan_core.security.bounded_process import run_bounded_operation
|
||||
from govoplan_dataflow.backend.backends.reference import ReferenceExecutionBackend
|
||||
from govoplan_dataflow.backend.executor import EXECUTOR_VERSION, PipelineExecutionError
|
||||
from govoplan_dataflow.backend.graph import validate_graph
|
||||
from govoplan_dataflow.backend.schemas import GraphEdge, GraphNode, GraphPosition, PipelineGraph
|
||||
from govoplan_dataflow.backend.service import _execute_pipeline_preview
|
||||
|
||||
|
||||
OUTER_ROWS = ({"id": "outer", "amount": 15}, {"id": "outer-low", "amount": 5})
|
||||
INNER_ROWS = ({"id": "inner", "amount": 25}, {"id": "inner-low", "amount": 2})
|
||||
|
||||
|
||||
def node(node_id: str, node_type: str, config: dict) -> GraphNode:
|
||||
return GraphNode(
|
||||
id=node_id, type=node_type, label=node_id,
|
||||
position=GraphPosition(x=0, y=0), config=config,
|
||||
)
|
||||
|
||||
|
||||
def nested_source_graph() -> PipelineGraph:
|
||||
# Both external sources deliberately share their node ID and logical name.
|
||||
# The pinned subflow still has exactly one distinct inline input binding.
|
||||
nested = PipelineGraph(
|
||||
nodes=[
|
||||
node("input", "source.inline", {
|
||||
"source_name": "bound_input", "rows": [], "input_binding": True,
|
||||
}),
|
||||
node("shared", "source.reference", {
|
||||
"source_name": "records", "source_ref": {"$parameter": "source_ref"},
|
||||
"expected_fingerprint": "inner-pinned", "consistency": "frozen",
|
||||
}),
|
||||
node("union", "combine.union", {"mode": "all"}),
|
||||
node("minimum", "filter.expression", {"expression": "amount >= ${minimum}"}),
|
||||
node("output", "output", {}),
|
||||
],
|
||||
edges=[
|
||||
GraphEdge(id="input-union", source="input", target="union"),
|
||||
GraphEdge(id="shared-union", source="shared", target="union"),
|
||||
GraphEdge(id="union-minimum", source="union", target="minimum"),
|
||||
GraphEdge(id="minimum-output", source="minimum", target="output"),
|
||||
],
|
||||
)
|
||||
return PipelineGraph(
|
||||
nodes=[
|
||||
node("shared", "source.reference", {
|
||||
"source_name": "records", "source_ref": "datasource:outer",
|
||||
"expected_fingerprint": "outer-pinned", "consistency": "current",
|
||||
}),
|
||||
node("nested", "subflow", {
|
||||
"template_ref": "fixture-nested-source", "template_version": "1",
|
||||
"parameters": {"source_ref": "datasource:inner", "minimum": 10},
|
||||
"graph": nested.model_dump(mode="python"),
|
||||
}),
|
||||
node("output", "output", {}),
|
||||
],
|
||||
edges=[
|
||||
GraphEdge(id="shared-nested", source="shared", target="nested"),
|
||||
GraphEdge(id="nested-output", source="nested", target="output"),
|
||||
],
|
||||
)
|
||||
|
||||
|
||||
class ReferenceSubflowProcessTests(unittest.TestCase):
|
||||
def setUp(self) -> None:
|
||||
self.graph = nested_source_graph()
|
||||
self.assertEqual([], [item.model_dump() for item in validate_graph(self.graph) if item.severity == "error"])
|
||||
self.session = Mock(spec=Session)
|
||||
self.principal = ApiPrincipal(
|
||||
principal=PrincipalRef(
|
||||
account_id="fixture-account", membership_id="fixture-membership",
|
||||
tenant_id="fixture-tenant", scopes=frozenset(),
|
||||
),
|
||||
account=object(), user=object(),
|
||||
)
|
||||
self.provider = Mock(spec=DatasourceCatalogueProvider)
|
||||
self.provider.read_datasource.side_effect = self.read_source
|
||||
self.registry = SimpleNamespace(
|
||||
has_capability=lambda name: name == CAPABILITY_DATASOURCE_CATALOGUE,
|
||||
capability=lambda name: self.provider if name == CAPABILITY_DATASOURCE_CATALOGUE else None,
|
||||
)
|
||||
self.reads: list[tuple[int, DatasourceReadRequest]] = []
|
||||
self.denied = False
|
||||
|
||||
def read_source(self, session, principal, *, request: DatasourceReadRequest) -> DatasourceReadResult:
|
||||
self.assertIs(session, self.session)
|
||||
self.assertIs(principal, self.principal)
|
||||
self.reads.append((os.getpid(), request))
|
||||
if request.datasource_ref == "datasource:inner" and self.denied:
|
||||
raise DatasourceAccessError("Current principal cannot read datasource:inner.")
|
||||
rows, fingerprint = {
|
||||
"datasource:outer": (OUTER_ROWS, "outer-pinned"),
|
||||
"datasource:inner": (INNER_ROWS, "inner-pinned"),
|
||||
}[request.datasource_ref]
|
||||
self.assertEqual(fingerprint, request.expected_fingerprint)
|
||||
return DatasourceReadResult(
|
||||
datasource=DatasourceDescriptor(
|
||||
ref=request.datasource_ref, source_name="records", name="Fixture source",
|
||||
kind="custom", mode="static", shape="tabular",
|
||||
fingerprint=fingerprint, provider="fixture-catalogue",
|
||||
),
|
||||
rows=rows, total_rows=len(rows), truncated=False,
|
||||
)
|
||||
|
||||
def preview(self, *, row_limit: int = 10):
|
||||
return _execute_pipeline_preview(
|
||||
self.graph, session=self.session, principal=self.principal,
|
||||
registry=self.registry, backend="reference", row_limit=row_limit,
|
||||
preview_node_id="nested",
|
||||
)
|
||||
|
||||
def test_nested_parameterized_sources_execute_in_real_child_without_id_aliasing(self) -> None:
|
||||
with (
|
||||
patch.object(ReferenceExecutionBackend, "_execute_in_process", side_effect=AssertionError("parent evaluation")),
|
||||
patch("govoplan_dataflow.backend.backends.reference.run_bounded_operation", wraps=run_bounded_operation) as worker,
|
||||
):
|
||||
result, version = self.preview()
|
||||
worker.assert_called_once()
|
||||
self.assertEqual(EXECUTOR_VERSION, version)
|
||||
self.assertEqual([OUTER_ROWS[0], INNER_ROWS[0]], result.rows)
|
||||
self.assertEqual(2, result.total_rows)
|
||||
self.assertFalse(result.truncated)
|
||||
self.assertEqual(2, result.input_row_count) # Root input count, not the nested binding again.
|
||||
self.assertIsNotNone(result.node_preview)
|
||||
self.assertEqual("nested", result.node_preview.node_id)
|
||||
self.assertEqual(result.rows, result.node_preview.rows)
|
||||
self.assertEqual(2, result.node_preview.total_rows)
|
||||
self.assertEqual(
|
||||
[("shared", "succeeded", 0, 2), ("nested", "succeeded", 2, 2), ("output", "succeeded", 2, 2)],
|
||||
[(item.node_id, item.status, item.input_rows, item.output_rows) for item in result.node_diagnostics],
|
||||
)
|
||||
binding_hash = hashlib.sha256(json.dumps(list(OUTER_ROWS), sort_keys=True, separators=(",", ":")).encode()).hexdigest()
|
||||
self.assertEqual([
|
||||
{"node_id": "shared", "source_ref": "datasource:outer", "source_name": "records", "kind": "datasource", "provider": "fixture-catalogue", "fingerprint": "outer-pinned", "row_count": 2, "preview_rows": 2, "truncated": False},
|
||||
{"node_id": "input", "source_name": "bound_input", "kind": "inline", "fingerprint": binding_hash, "row_count": 2, "subflow_node_id": "nested"},
|
||||
{"node_id": "shared", "source_ref": "datasource:inner", "source_name": "records", "kind": "datasource", "provider": "fixture-catalogue", "fingerprint": "inner-pinned", "row_count": 2, "preview_rows": 2, "truncated": False, "subflow_node_id": "nested"},
|
||||
], result.source_fingerprints)
|
||||
self.assertEqual(
|
||||
[("datasource:outer", "current", "outer-pinned"), ("datasource:inner", "frozen", "inner-pinned")],
|
||||
[(request.datasource_ref, request.consistency, request.expected_fingerprint) for _, request in self.reads],
|
||||
)
|
||||
self.assertTrue(all(pid == os.getpid() for pid, _ in self.reads))
|
||||
self.assertTrue(all(request.limit <= 500 and request.offset == 0 for _, request in self.reads))
|
||||
self.assertEqual([], self.session.mock_calls)
|
||||
|
||||
def test_nested_result_and_node_preview_keep_full_totals_when_output_is_bounded(self) -> None:
|
||||
result, _ = self.preview(row_limit=1)
|
||||
self.assertEqual([OUTER_ROWS[0]], result.rows)
|
||||
self.assertEqual(2, result.total_rows)
|
||||
self.assertTrue(result.truncated)
|
||||
self.assertEqual([OUTER_ROWS[0]], result.node_preview.rows)
|
||||
self.assertEqual(2, result.node_preview.total_rows)
|
||||
self.assertTrue(result.node_preview.truncated)
|
||||
self.assertEqual({"outer-pinned", "inner-pinned"}, {
|
||||
item["fingerprint"] for item in result.source_fingerprints if item["kind"] == "datasource"
|
||||
})
|
||||
self.assertEqual([], self.session.mock_calls)
|
||||
|
||||
def test_denied_nested_datasource_stops_before_worker_or_persistence(self) -> None:
|
||||
self.denied = True
|
||||
with (
|
||||
patch("govoplan_dataflow.backend.service.execute_typed_graph") as execute,
|
||||
patch("govoplan_dataflow.backend.backends.reference.run_bounded_operation") as worker,
|
||||
self.assertRaises(PipelineExecutionError) as caught,
|
||||
):
|
||||
self.preview()
|
||||
execute.assert_not_called()
|
||||
worker.assert_not_called()
|
||||
self.assertEqual("shared", caught.exception.node_id)
|
||||
self.assertEqual("Current principal cannot read datasource:inner.", str(caught.exception))
|
||||
self.assertFalse(caught.exception.retryable)
|
||||
self.assertIsInstance(caught.exception.__cause__, DatasourceAccessError)
|
||||
self.assertEqual(["datasource:outer", "datasource:inner"], [request.datasource_ref for _, request in self.reads])
|
||||
self.assertEqual([], self.session.mock_calls)
|
||||
+4
-3
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@govoplan/dataflow-webui",
|
||||
"version": "0.1.23",
|
||||
"version": "0.1.25",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"main": "src/index.ts",
|
||||
@@ -15,10 +15,11 @@
|
||||
},
|
||||
"scripts": {
|
||||
"typecheck": "tsc --noEmit",
|
||||
"test:structure": "node scripts/test-dataflow-page-structure.mjs"
|
||||
"test:structure": "node scripts/test-dataflow-page-structure.mjs",
|
||||
"test:save-completion": "node --test scripts/test-save-completion.mjs"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"@govoplan/core-webui": "^0.1.18",
|
||||
"@govoplan/core-webui": "^0.1.46",
|
||||
"@xyflow/react": "^12.11.2",
|
||||
"lucide-react": "^1.23.0",
|
||||
"react": ">=19.2.7 <20",
|
||||
|
||||
Executable
+200
@@ -0,0 +1,200 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { readFileSync } from "node:fs";
|
||||
import { createRequire } from "node:module";
|
||||
import test from "node:test";
|
||||
import vm from "node:vm";
|
||||
|
||||
const require = createRequire(new URL("../../../govoplan-core/webui/package.json", import.meta.url));
|
||||
const { transformSync } = require("esbuild");
|
||||
const page = readFileSync(new URL("../src/features/dataflow/DataflowPage.tsx", import.meta.url), "utf8");
|
||||
function loadTs(path) {
|
||||
const context = vm.createContext({ module: { exports: {} }, require: () => ({}), structuredClone });
|
||||
context.exports = context.module.exports;
|
||||
vm.runInContext(transformSync(readFileSync(new URL(path, import.meta.url), "utf8"), { loader: "ts", format: "cjs" }).code, context);
|
||||
return context.module.exports;
|
||||
}
|
||||
const { reconcilePipelineSave } = loadTs("../src/features/dataflow/saveCompletion.ts");
|
||||
const { draftFingerprint, pipelinePayload } = loadTs("../src/features/dataflow/model.ts");
|
||||
const { authAuthorityKey } = loadTs("../../../govoplan-core/webui/src/api/authAuthority.ts");
|
||||
const start = page.indexOf(" const saveDraft = useCallback(async (): Promise<boolean> => {");
|
||||
const end = page.indexOf(" }, [canEdit, draft, settings, authorityKey, authorityGeneration]);", start);
|
||||
assert.ok(start >= 0 && end > start, "exercise the actual page save closure");
|
||||
const code = transformSync(page.slice(start, end) + " }, []);\nmodule.exports = saveDraft;", { loader: "ts" }).code;
|
||||
const base = () => ({
|
||||
id: "pipeline-1", currentRevision: 1, name: "Pipeline", description: "Submitted", status: "draft",
|
||||
graph: { nodes: [{ id: "a", config: { rows: [{ id: "001", value: " text " }] } }, { id: "b" }], edges: [{ id: "a-b" }] },
|
||||
sqlText: "", editorMode: "graph", scopeType: "tenant", scopeId: "tenant-1", definitionKind: "flow",
|
||||
inheritToLowerScopes: false, allowRun: true, allowReuse: true, allowAutomation: false, governance: { actions: {} }
|
||||
});
|
||||
function harness(draft = base()) {
|
||||
const calls = [];
|
||||
const responses = [];
|
||||
const context = vm.createContext({
|
||||
module: { exports: {} }, useCallback: (value) => value, structuredClone,
|
||||
draft, canEdit: true, settings: {}, saveInFlight: { current: false },
|
||||
draftSession: { current: { generation: 1, value: draft } },
|
||||
authorityKey: "authority-A", saveContext: { current: "authority-A" }, unresolvedSaveGeneration: { current: null },
|
||||
authorityGeneration: 0, saveAuthorityEpoch: { current: { key: "authority-A", revision: 0 } },
|
||||
reconcilePipelineSave, draftFingerprint, pipelinePayload, draftFromPipeline: (value) => value,
|
||||
updateDataflowPipeline: (_settings, id, payload) => new Promise((resolve, reject) => { calls.push({ id, payload }); responses.push({ resolve, reject }); }),
|
||||
createDataflowPipeline: (_settings, payload) => new Promise((resolve, reject) => { calls.push({ payload }); responses.push({ resolve, reject }); }),
|
||||
setDraftValue: (value) => { context.draft = value; },
|
||||
setSavedDraft: (value) => { context.baseline = value; },
|
||||
setSaving: (value) => { context.saving = value; },
|
||||
setError: (value) => { context.error = value; },
|
||||
setSuccess: () => {}, setPipelines: () => {}, setSelectedNodeId: () => {}, setDiagnostics: () => {}, apiErrorMessage: String
|
||||
});
|
||||
vm.runInContext(code, context);
|
||||
return { context, calls, responses, save: context.module.exports };
|
||||
}
|
||||
function accepted(draft, revision = 2) {
|
||||
return { ...structuredClone(draft), currentRevision: revision, current_revision: revision, governance: { actions: { edit: { allowed: true } } } };
|
||||
}
|
||||
|
||||
test("edits made during an accepted save survive and navigation remains blocked until saved", async () => {
|
||||
const h = harness();
|
||||
const submitted = h.context.draft;
|
||||
const pending = h.save();
|
||||
const newerGraph = { ...submitted.graph, nodes: [...submitted.graph.nodes].reverse(), custom: { preserve: ["001", null, ""] } };
|
||||
h.context.draftSession.current.value = { ...submitted, description: "Typed during save", graph: newerGraph };
|
||||
h.responses[0].resolve(accepted(submitted));
|
||||
assert.equal(await pending, false);
|
||||
assert.equal(h.context.draft.description, "Typed during save");
|
||||
assert.equal(h.context.draft.graph, newerGraph, "graph remains atomic, including order and unknown data");
|
||||
assert.equal(h.context.baseline.description, "Submitted");
|
||||
assert.equal(h.context.draft.currentRevision, 2);
|
||||
assert.notEqual(draftFingerprint(h.context.draft), draftFingerprint(h.context.baseline));
|
||||
const second = h.save();
|
||||
assert.equal(h.calls[1].payload.expected_revision, 2, "next save uses the accepted revision, not the stale submitted one");
|
||||
assert.deepEqual(h.calls[1].payload.graph, newerGraph);
|
||||
h.responses[1].resolve(accepted(h.context.draft, 3));
|
||||
assert.equal(await second, true);
|
||||
assert.equal(draftFingerprint(h.context.draft), draftFingerprint(h.context.baseline));
|
||||
});
|
||||
|
||||
test("unchanged submitted fields accept canonical response and new identities without a duplicate create", async () => {
|
||||
const h = harness({ ...base(), id: null, currentRevision: null });
|
||||
const pending = h.save();
|
||||
assert.equal(await h.save(), false);
|
||||
assert.equal(h.calls.length, 1);
|
||||
const response = { ...accepted(h.context.draft), id: "created", name: "Server canonical name" };
|
||||
h.responses[0].resolve(response);
|
||||
assert.equal(await pending, true);
|
||||
assert.equal(h.context.draft.id, "created");
|
||||
assert.equal(h.context.draft.name, "Server canonical name");
|
||||
});
|
||||
|
||||
test("replacement draft, authority change and unmount cannot receive an old save completion", async () => {
|
||||
for (const change of ["replacement", "authority", "unmount"]) {
|
||||
const h = harness();
|
||||
const pending = h.save();
|
||||
if (change === "authority") h.context.saveContext.current = "authority-B";
|
||||
else h.context.draftSession.current.generation += 1;
|
||||
h.responses[0].resolve(accepted(h.context.draft));
|
||||
assert.equal(await pending, false);
|
||||
assert.equal(h.context.baseline, undefined);
|
||||
assert.equal(h.context.draft.currentRevision, 1);
|
||||
}
|
||||
});
|
||||
|
||||
test("harmless session/profile object refresh preserves an accepted new identity and revision", async () => {
|
||||
const h = harness({ ...base(), id: null, currentRevision: null });
|
||||
const settings = { apiBaseUrl: "/api", apiKey: "", accessToken: "" };
|
||||
const auth = {
|
||||
user: { id: "member", account_id: "account", email: "person@example.test", display_name: "Before" },
|
||||
tenant: { id: "tenant" }, scopes: ["dataflow:pipeline:write"], roles: [], groups: []
|
||||
};
|
||||
h.context.authorityKey = authAuthorityKey(auth, settings);
|
||||
h.context.saveContext.current = h.context.authorityKey;
|
||||
const pending = h.save();
|
||||
h.context.saveContext.current = authAuthorityKey({ ...structuredClone(auth), user: { ...auth.user, display_name: "After", preferred_language: "de" } }, { ...settings });
|
||||
h.responses[0].resolve({ ...accepted(h.context.draft), id: "accepted-created-id" });
|
||||
assert.equal(await pending, true);
|
||||
assert.equal(h.context.draft.id, "accepted-created-id");
|
||||
const next = h.save();
|
||||
assert.equal(h.calls[1].id, "accepted-created-id", "retry updates the accepted identity, never creates another pipeline");
|
||||
assert.equal(h.calls[1].payload.expected_revision, 2);
|
||||
h.responses[1].resolve(accepted(h.context.draft, 3));
|
||||
assert.equal(await next, true);
|
||||
});
|
||||
|
||||
test("an accepted save across a real authority change cannot be blindly retried as a duplicate create", async () => {
|
||||
const h = harness({ ...base(), id: null, currentRevision: null });
|
||||
const pending = h.save();
|
||||
h.context.saveContext.current = "authority-B";
|
||||
h.responses[0].resolve({ ...accepted(h.context.draft), id: "accepted-under-A" });
|
||||
assert.equal(await pending, false);
|
||||
h.context.authorityKey = "authority-B";
|
||||
assert.equal(await h.save(), false);
|
||||
assert.equal(h.calls.length, 1);
|
||||
assert.match(h.context.error, /Reload and review/);
|
||||
});
|
||||
|
||||
test("returning to authority A after B does not revive a stale A save completion", async () => {
|
||||
const h = harness();
|
||||
const pending = h.save();
|
||||
h.context.saveAuthorityEpoch.current = { key: "authority-A", revision: 2 };
|
||||
h.responses[0].resolve(accepted(h.context.draft));
|
||||
assert.equal(await pending, false);
|
||||
assert.equal(h.context.baseline, undefined);
|
||||
});
|
||||
|
||||
test("conflict preserves both local draft and prior revision, allowing an explicit reviewed retry", async () => {
|
||||
const h = harness();
|
||||
const original = h.context.draft;
|
||||
const pending = h.save();
|
||||
h.responses[0].reject(new Error("revision conflict"));
|
||||
assert.equal(await pending, false);
|
||||
assert.equal(h.context.draft, original);
|
||||
assert.equal(h.context.baseline, undefined);
|
||||
assert.match(h.context.error, /revision conflict/);
|
||||
assert.equal(h.context.saveInFlight.current, false);
|
||||
});
|
||||
|
||||
test("the submitted baseline is frozen even if a nested local editor mutates a shared object", async () => {
|
||||
const h = harness();
|
||||
const serverAccepted = accepted(h.context.draft);
|
||||
const pending = h.save();
|
||||
h.context.draft.graph.nodes.reverse();
|
||||
h.responses[0].resolve(serverAccepted);
|
||||
assert.equal(await pending, false);
|
||||
assert.equal(h.context.draft.graph.nodes[0].id, "b");
|
||||
assert.equal(h.context.baseline.graph.nodes[0].id, "a");
|
||||
assert.equal(h.calls[0].payload.graph.nodes[0].id, "a", "submission data is not aliased to later editor mutations");
|
||||
});
|
||||
|
||||
test("CSV imports explicitly preserve text by default; JSON payload stays independent", () => {
|
||||
assert.match(page, /\[csvValueMode, setCsvValueMode\] = useState<"text" \| "legacy_typed">\("text"\)/);
|
||||
assert.match(page, /\? \{ format, rows \}\s*: \{ format, csv_text: csvText, delimiter, csv_value_mode: csvValueMode \}/);
|
||||
assert.match(page, /setCsvValueMode\("text"\)/);
|
||||
});
|
||||
|
||||
test("source dialog sends exact CSV content and selected mode without changing JSON rows", async () => {
|
||||
const dialog = page.slice(page.indexOf("function SourceSnapshotDialog("));
|
||||
const start = dialog.indexOf(" const create = async (): Promise<boolean> => {");
|
||||
const end = dialog.indexOf("\n };", start);
|
||||
assert.ok(start >= 0 && end > start);
|
||||
const code = transformSync(dialog.slice(start, end) + "\n}; module.exports = create;", { loader: "ts" }).code;
|
||||
const csvText = 'code,value\r\n001," text "\r\n';
|
||||
for (const format of ["csv", "json"]) {
|
||||
for (const csvValueMode of ["text", "legacy_typed"]) {
|
||||
let payload;
|
||||
const context = vm.createContext({
|
||||
module: { exports: {} }, settings: {}, format, csvValueMode, csvText, delimiter: ",",
|
||||
name: "Fixture", sourceName: "fixture", description: "", rowsText: '[{"code":"001","value":" text "}]',
|
||||
isRecord: (value) => Boolean(value) && typeof value === "object" && !Array.isArray(value),
|
||||
createDataflowSourceSnapshot: async (_settings, value) => { payload = value; return {}; },
|
||||
onCreated: () => {}, setBusy: () => {}, setError: () => {}, apiErrorMessage: String
|
||||
});
|
||||
vm.runInContext(code, context);
|
||||
assert.equal(await context.module.exports(), true);
|
||||
if (format === "csv") {
|
||||
assert.equal(payload.csv_value_mode, csvValueMode);
|
||||
assert.equal(payload.csv_text, csvText);
|
||||
} else {
|
||||
assert.equal("csv_value_mode" in payload, false);
|
||||
assert.equal(JSON.stringify(payload.rows), '[{"code":"001","value":" text "}]');
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
@@ -392,6 +392,7 @@ export function createDataflowSourceSnapshot(
|
||||
format: "csv";
|
||||
csv_text: string;
|
||||
delimiter: string;
|
||||
csv_value_mode?: "text" | "legacy_typed";
|
||||
}
|
||||
)
|
||||
): Promise<TabularSource> {
|
||||
|
||||
@@ -56,6 +56,7 @@ import { DialogSection, ActionToolbar,
|
||||
WorkspaceFrame,
|
||||
WorkspaceLayout,
|
||||
hasScope,
|
||||
authAuthorityKey,
|
||||
isApiError,
|
||||
useUnsavedChanges,
|
||||
useUnsavedDraftGuard,
|
||||
@@ -129,6 +130,8 @@ import {
|
||||
DATAFLOW_RUN_DOCUMENTATION
|
||||
} from "./interfacePatterns";
|
||||
|
||||
import { reconcilePipelineSave } from "./saveCompletion";
|
||||
|
||||
type ResultTab = "preview" | "diagnostics";
|
||||
type SnapshotFormat = "json" | "csv";
|
||||
|
||||
@@ -143,7 +146,29 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings
|
||||
[location.search]
|
||||
);
|
||||
const [pipelines, setPipelines] = useState<Pipeline[]>([]);
|
||||
const [draft, setDraft] = useState<PipelineDraft | null>(null);
|
||||
const [draft, setDraftValue] = useState<PipelineDraft | null>(null);
|
||||
const draftSession = useRef<{ generation: number; value: PipelineDraft | null }>({ generation: 0, value: null });
|
||||
const saveInFlight = useRef(false);
|
||||
const authorityKey = authAuthorityKey(auth, settings);
|
||||
const saveAuthorityEpoch = useRef({ key: authorityKey, revision: 0 });
|
||||
if (saveAuthorityEpoch.current.key !== authorityKey) {
|
||||
saveAuthorityEpoch.current = { key: authorityKey, revision: saveAuthorityEpoch.current.revision + 1 };
|
||||
}
|
||||
const authorityGeneration = saveAuthorityEpoch.current.revision;
|
||||
const saveContext = useRef(authorityKey);
|
||||
saveContext.current = authorityKey;
|
||||
const unresolvedSaveGeneration = useRef<number | null>(null);
|
||||
useEffect(() => {
|
||||
saveContext.current = authorityKey;
|
||||
return () => { if (saveContext.current === authorityKey) saveContext.current = ""; };
|
||||
}, [authorityKey]);
|
||||
// Replacement (selection, reload, discard, derive) is a different edit session,
|
||||
// even when both unsaved drafts have a null identifier.
|
||||
const setDraft = useCallback((next: PipelineDraft | null) => {
|
||||
draftSession.current = { generation: draftSession.current.generation + 1, value: next };
|
||||
setDraftValue(next);
|
||||
}, []);
|
||||
useEffect(() => () => { draftSession.current.generation += 1; }, []);
|
||||
const [savedDraft, setSavedDraft] = useState<PipelineDraft | null>(null);
|
||||
const [selectedNodeId, setSelectedNodeId] = useState<string | null>(null);
|
||||
const [search, setSearch] = useState("");
|
||||
@@ -322,15 +347,27 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings
|
||||
}, [savedDraft]);
|
||||
|
||||
const saveDraft = useCallback(async (): Promise<boolean> => {
|
||||
if (saveInFlight.current) return false;
|
||||
if (authorityKey !== saveContext.current || authorityGeneration !== saveAuthorityEpoch.current.revision) return false;
|
||||
if (unresolvedSaveGeneration.current === draftSession.current.generation) {
|
||||
setError("A prior save completed after authorization changed. Reload and review the server revision before saving again.");
|
||||
return false;
|
||||
}
|
||||
if (!draft || !canEdit || !draft.name.trim()) {
|
||||
setError(!draft?.name.trim() ? "Pipeline name is required." : "You cannot save this pipeline.");
|
||||
return false;
|
||||
}
|
||||
saveInFlight.current = true;
|
||||
const generation = draftSession.current.generation;
|
||||
const context = authorityKey;
|
||||
const isCurrent = () => generation === draftSession.current.generation
|
||||
&& context === saveContext.current && authorityGeneration === saveAuthorityEpoch.current.revision;
|
||||
setSaving(true);
|
||||
setError("");
|
||||
setSuccess("");
|
||||
try {
|
||||
const payload = pipelinePayload(draft);
|
||||
const submitted = structuredClone(draft);
|
||||
const payload = pipelinePayload(submitted);
|
||||
const saved = draft.id && draft.currentRevision
|
||||
? await updateDataflowPipeline(settings, draft.id, {
|
||||
...payload,
|
||||
@@ -338,22 +375,32 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings
|
||||
})
|
||||
: await createDataflowPipeline(settings, payload);
|
||||
const next = draftFromPipeline(saved);
|
||||
setDraft(next);
|
||||
if (!isCurrent() || !draftSession.current.value) {
|
||||
if (generation === draftSession.current.generation) unresolvedSaveGeneration.current = generation;
|
||||
return false;
|
||||
}
|
||||
const reconciled = reconcilePipelineSave(submitted, draftSession.current.value, next);
|
||||
draftSession.current.value = reconciled;
|
||||
setDraftValue(reconciled);
|
||||
setSavedDraft(structuredClone(next));
|
||||
setPipelines((current) => [saved, ...current.filter((item) => item.id !== saved.id)]);
|
||||
setSelectedNodeId((current) => current && next.graph.nodes.some((node) => node.id === current)
|
||||
setSelectedNodeId((current) => current && reconciled.graph.nodes.some((node) => node.id === current)
|
||||
? current
|
||||
: next.graph.nodes[0]?.id ?? null);
|
||||
: reconciled.graph.nodes[0]?.id ?? null);
|
||||
setDiagnostics([]);
|
||||
setSuccess(`Saved revision ${saved.current_revision}.`);
|
||||
return true;
|
||||
const fullySaved = draftFingerprint(reconciled) === draftFingerprint(next);
|
||||
setSuccess(fullySaved ? `Saved revision ${saved.current_revision}.`
|
||||
: "The submitted revision was saved. Newer edits remain unsaved.");
|
||||
// A navigation guard may proceed only if ALL current edits were accepted.
|
||||
return fullySaved;
|
||||
} catch (saveError) {
|
||||
setError(apiErrorMessage(saveError));
|
||||
if (isCurrent()) setError(apiErrorMessage(saveError));
|
||||
return false;
|
||||
} finally {
|
||||
saveInFlight.current = false;
|
||||
setSaving(false);
|
||||
}
|
||||
}, [canEdit, draft, settings]);
|
||||
}, [canEdit, draft, settings, authorityKey, authorityGeneration]);
|
||||
|
||||
useUnsavedDraftGuard({
|
||||
dirty,
|
||||
@@ -402,7 +449,12 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings
|
||||
};
|
||||
|
||||
const updateDraft = (patch: Partial<PipelineDraft>) => {
|
||||
setDraft((current) => current ? { ...current, ...patch } : current);
|
||||
const current = draftSession.current.value;
|
||||
if (current) {
|
||||
const next = { ...current, ...patch };
|
||||
draftSession.current.value = next;
|
||||
setDraftValue(next);
|
||||
}
|
||||
setSuccess("");
|
||||
};
|
||||
|
||||
@@ -615,6 +667,22 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings
|
||||
|
||||
return (
|
||||
<WorkspaceFrame as="main" height="viewport" surface="plain" className="dataflow-page" label="Dataflow workspace">
|
||||
<WorkspaceActionBar
|
||||
title="Pipelines"
|
||||
titleHelp={<DocumentationHelpLink reference={DATAFLOW_DOCUMENTATION} />}
|
||||
scope="workspace"
|
||||
variant="collection"
|
||||
refreshable
|
||||
reloadAction={{ onReload: () => void loadPipelines(draft?.id), loading, label: "Refresh pipelines" }}
|
||||
createAction={<IconButton
|
||||
label="New pipeline"
|
||||
icon={<Plus size={17} />}
|
||||
variant="primary"
|
||||
onClick={createNew}
|
||||
disabled={!canWrite}
|
||||
disabledReason={!canWrite ? DATAFLOW_I18N.writeReason : undefined}
|
||||
/>}
|
||||
/>
|
||||
<WorkspaceLayout
|
||||
variant="split"
|
||||
primarySize="compact"
|
||||
@@ -625,21 +693,6 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings
|
||||
contentLabel="Pipeline editor"
|
||||
contentClassName="dataflow-workspace"
|
||||
primary={<>
|
||||
<WorkspaceActionBar
|
||||
scope="collection-pane"
|
||||
variant="collection"
|
||||
refreshable
|
||||
reloadAction={{ onReload: () => void loadPipelines(draft?.id), loading, label: "Refresh pipelines" }}
|
||||
contextActions={<strong>Pipelines</strong>}
|
||||
createAction={<IconButton
|
||||
label="New pipeline"
|
||||
icon={<Plus size={17} />}
|
||||
variant="primary"
|
||||
onClick={createNew}
|
||||
disabled={!canWrite}
|
||||
disabledReason={!canWrite ? DATAFLOW_I18N.writeReason : undefined}
|
||||
/>}
|
||||
/>
|
||||
<FilterBar surface="panel">
|
||||
<input
|
||||
type="search"
|
||||
@@ -707,7 +760,6 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings
|
||||
<option value="archived">Archived</option>
|
||||
</select>
|
||||
</div>}
|
||||
helpAction={<DocumentationHelpLink reference={DATAFLOW_DOCUMENTATION} />}
|
||||
primaryActions={<div className="dataflow-command-bar">
|
||||
<SegmentedControl<EditorMode>
|
||||
ariaLabel="Pipeline editor mode"
|
||||
@@ -1901,6 +1953,7 @@ function DataflowTriggersDialog({
|
||||
<Dialog
|
||||
open={open}
|
||||
title={`Automation · ${pipeline?.name ?? "pipeline"}`}
|
||||
titleHelp={<DocumentationHelpLink reference={DATAFLOW_FIELDS_DOCUMENTATION} />}
|
||||
className="dataflow-triggers-dialog"
|
||||
closeDisabled={busy}
|
||||
onClose={close}
|
||||
@@ -1957,7 +2010,6 @@ function DataflowTriggersDialog({
|
||||
{!busy && !triggers.length ? <small>No triggers configured</small> : null}
|
||||
</div>
|
||||
<div className="dataflow-trigger-form">
|
||||
<DocumentationHelpLink reference={DATAFLOW_FIELDS_DOCUMENTATION} />
|
||||
{error ? <DismissibleAlert tone="danger" resetKey={error}>{error}</DismissibleAlert> : null}
|
||||
<FormField label="Name" documentation={DATAFLOW_FIELDS_DOCUMENTATION}>
|
||||
<input disabled={!editable} value={name} onChange={(event) => setName(event.target.value)} />
|
||||
@@ -2261,6 +2313,7 @@ function RunPipelineDialog({
|
||||
<Dialog
|
||||
open={open}
|
||||
title={`Run ${pipeline?.name ?? "pipeline"}`}
|
||||
titleHelp={<DocumentationHelpLink reference={DATAFLOW_RUN_DOCUMENTATION} />}
|
||||
className="dataflow-run-dialog"
|
||||
onClose={() => {
|
||||
if (!busy) {
|
||||
@@ -2297,7 +2350,6 @@ function RunPipelineDialog({
|
||||
</DismissibleAlert>
|
||||
) : null}
|
||||
<div className="dataflow-run-controls">
|
||||
<DocumentationHelpLink reference={DATAFLOW_RUN_DOCUMENTATION} />
|
||||
<SegmentedControl<RunMode>
|
||||
ariaLabel="Run output"
|
||||
options={[
|
||||
@@ -2518,6 +2570,7 @@ function SourceSnapshotDialog({
|
||||
const [rowsText, setRowsText] = useState("[]");
|
||||
const [csvText, setCsvText] = useState("");
|
||||
const [delimiter, setDelimiter] = useState(",");
|
||||
const [csvValueMode, setCsvValueMode] = useState<"text" | "legacy_typed">("text");
|
||||
const [fileInputKey, setFileInputKey] = useState(0);
|
||||
const [busy, setBusy] = useState(false);
|
||||
const [error, setError] = useState("");
|
||||
@@ -2531,6 +2584,7 @@ function SourceSnapshotDialog({
|
||||
|| rowsText !== "[]"
|
||||
|| csvText !== ""
|
||||
|| delimiter !== ","
|
||||
|| csvValueMode !== "text"
|
||||
)
|
||||
);
|
||||
|
||||
@@ -2542,6 +2596,7 @@ function SourceSnapshotDialog({
|
||||
setRowsText("[]");
|
||||
setCsvText("");
|
||||
setDelimiter(",");
|
||||
setCsvValueMode("text");
|
||||
setFileInputKey((current) => current + 1);
|
||||
setError("");
|
||||
};
|
||||
@@ -2577,7 +2632,7 @@ function SourceSnapshotDialog({
|
||||
description: description.trim() || null,
|
||||
...(format === "json"
|
||||
? { format, rows }
|
||||
: { format, csv_text: csvText, delimiter })
|
||||
: { format, csv_text: csvText, delimiter, csv_value_mode: csvValueMode })
|
||||
});
|
||||
onCreated(source);
|
||||
return true;
|
||||
@@ -2682,6 +2737,16 @@ function SourceSnapshotDialog({
|
||||
<option value="|">Pipe</option>
|
||||
</select>
|
||||
</FormField>
|
||||
<FormField label="CSV values" documentation={DATAFLOW_FIELDS_DOCUMENTATION}>
|
||||
<select
|
||||
value={csvValueMode}
|
||||
onChange={(event) => setCsvValueMode(event.target.value as "text" | "legacy_typed")}
|
||||
disabled={busy}
|
||||
>
|
||||
<option value="text">Preserve text (no automatic conversion)</option>
|
||||
<option value="legacy_typed">Infer types (legacy)</option>
|
||||
</select>
|
||||
</FormField>
|
||||
<FormField label="CSV data" documentation={DATAFLOW_FIELDS_DOCUMENTATION}>
|
||||
<textarea
|
||||
className="dataflow-json-editor"
|
||||
|
||||
+18
@@ -0,0 +1,18 @@
|
||||
import type { PipelineDraft } from "./model";
|
||||
|
||||
/** Reconcile one accepted save, never structurally merge/reorder graph data.
|
||||
* Fields edited since submission stay local; identity, revision and authority
|
||||
* always come from the accepted server response.
|
||||
*/
|
||||
export function reconcilePipelineSave(
|
||||
submitted: PipelineDraft, current: PipelineDraft, accepted: PipelineDraft
|
||||
): PipelineDraft {
|
||||
const result = { ...current, ...accepted };
|
||||
for (const key of Object.keys(submitted) as Array<keyof PipelineDraft>) {
|
||||
if (key === "id" || key === "currentRevision" || key === "governance") continue;
|
||||
if (JSON.stringify(current[key]) !== JSON.stringify(submitted[key])) {
|
||||
Object.assign(result, { [key]: current[key] });
|
||||
}
|
||||
}
|
||||
return result;
|
||||
}
|
||||
@@ -1,6 +1,11 @@
|
||||
import type { PlatformTranslations } from "@govoplan/core-webui";
|
||||
|
||||
const en = {
|
||||
"A prior save completed after authorization changed. Reload and review the server revision before saving again.": "A prior save completed after authorization changed. Reload and review the server revision before saving again.",
|
||||
"The submitted revision was saved. Newer edits remain unsaved.": "The submitted revision was saved. Newer edits remain unsaved.",
|
||||
"CSV values": "CSV values",
|
||||
"Preserve text (no automatic conversion)": "Preserve text (no automatic conversion)",
|
||||
"Infer types (legacy)": "Infer types (legacy)",
|
||||
"i18n:govoplan-dataflow.dataflow": "Dataflow",
|
||||
"i18n:govoplan-dataflow.library": "Pipeline library",
|
||||
"i18n:govoplan-dataflow.graph": "Graph editor",
|
||||
@@ -100,6 +105,11 @@ const en = {
|
||||
} as const;
|
||||
|
||||
const de: Record<keyof typeof en, string> = {
|
||||
"A prior save completed after authorization changed. Reload and review the server revision before saving again.": "Ein vorheriger Speichervorgang wurde nach einer Berechtigungsänderung abgeschlossen. Vor erneutem Speichern neu laden und die Serverrevision prüfen.",
|
||||
"The submitted revision was saved. Newer edits remain unsaved.": "Die übermittelte Revision wurde gespeichert. Neuere Änderungen sind noch ungespeichert.",
|
||||
"CSV values": "CSV-Werte",
|
||||
"Preserve text (no automatic conversion)": "Text erhalten (keine automatische Umwandlung)",
|
||||
"Infer types (legacy)": "Typen ableiten (bisheriges Verhalten)",
|
||||
"i18n:govoplan-dataflow.dataflow": "Datenfluss",
|
||||
"i18n:govoplan-dataflow.library": "Datenflussbibliothek",
|
||||
"i18n:govoplan-dataflow.graph": "Graph-Editor",
|
||||
|
||||
Reference in New Issue
Block a user