diff --git a/package.json b/package.json index 4656574..32e3fbe 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@govoplan/dataflow-webui", - "version": "0.1.24", + "version": "0.1.25", "private": true, "type": "module", "main": "webui/src/index.ts", @@ -14,7 +14,7 @@ "./styles/dataflow.css": "./webui/src/styles/dataflow.css" }, "peerDependencies": { - "@govoplan/core-webui": "^0.1.45", + "@govoplan/core-webui": "^0.1.46", "@xyflow/react": "^12.11.2", "lucide-react": "^1.23.0", "react": ">=19.2.7 <20", diff --git a/pyproject.toml b/pyproject.toml index d359928..b44babd 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,14 +4,14 @@ build-backend = "setuptools.build_meta" [project] name = "govoplan-dataflow" -version = "0.1.24" +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.45", + "govoplan-core>=0.1.46", "sqlglot>=30.14,<31", ] diff --git a/src/govoplan_dataflow/__init__.py b/src/govoplan_dataflow/__init__.py index 2128b18..db3fcf6 100644 --- a/src/govoplan_dataflow/__init__.py +++ b/src/govoplan_dataflow/__init__.py @@ -1,3 +1,3 @@ from __future__ import annotations -__version__ = "0.1.24" +__version__ = "0.1.25" diff --git a/src/govoplan_dataflow/backend/german_documentation.py b/src/govoplan_dataflow/backend/german_documentation.py index 10e0133..0eebbbc 100644 --- a/src/govoplan_dataflow/backend/german_documentation.py +++ b/src/govoplan_dataflow/backend/german_documentation.py @@ -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 = { @@ -56,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) diff --git a/src/govoplan_dataflow/backend/manifest.py b/src/govoplan_dataflow/backend/manifest.py index 15650b6..e689a7d 100644 --- a/src/govoplan_dataflow/backend/manifest.py +++ b/src/govoplan_dataflow/backend/manifest.py @@ -63,7 +63,7 @@ from govoplan_dataflow.backend.german_documentation import ( MODULE_ID = "dataflow" MODULE_NAME = "Dataflow" -MODULE_VERSION = "0.1.24" +MODULE_VERSION = "0.1.25" READ_SCOPE = "dataflow:pipeline:read" WRITE_SCOPE = "dataflow:pipeline:write" @@ -154,6 +154,64 @@ 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", diff --git a/src/govoplan_dataflow/backend/router.py b/src/govoplan_dataflow/backend/router.py index a9e27b0..29474ed 100644 --- a/src/govoplan_dataflow/backend/router.py +++ b/src/govoplan_dataflow/backend/router.py @@ -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( diff --git a/src/govoplan_dataflow/backend/schemas.py b/src/govoplan_dataflow/backend/schemas.py index 1fda064..1410849 100644 --- a/src/govoplan_dataflow/backend/schemas.py +++ b/src/govoplan_dataflow/backend/schemas.py @@ -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") diff --git a/tests/test_csv_staging_route.py b/tests/test_csv_staging_route.py new file mode 100755 index 0000000..0891aaf --- /dev/null +++ b/tests/test_csv_staging_route.py @@ -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() diff --git a/webui/package.json b/webui/package.json index bd26623..be42258 100644 --- a/webui/package.json +++ b/webui/package.json @@ -1,6 +1,6 @@ { "name": "@govoplan/dataflow-webui", - "version": "0.1.24", + "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.45", + "@govoplan/core-webui": "^0.1.46", "@xyflow/react": "^12.11.2", "lucide-react": "^1.23.0", "react": ">=19.2.7 <20", diff --git a/webui/scripts/test-save-completion.mjs b/webui/scripts/test-save-completion.mjs new file mode 100755 index 0000000..f7f81d2 --- /dev/null +++ b/webui/scripts/test-save-completion.mjs @@ -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 => {"); +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 => {"); + 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 "}]'); + } + } + } +}); diff --git a/webui/src/api/dataflow.ts b/webui/src/api/dataflow.ts index 8c37377..de808a1 100644 --- a/webui/src/api/dataflow.ts +++ b/webui/src/api/dataflow.ts @@ -392,6 +392,7 @@ export function createDataflowSourceSnapshot( format: "csv"; csv_text: string; delimiter: string; + csv_value_mode?: "text" | "legacy_typed"; } ) ): Promise { diff --git a/webui/src/features/dataflow/DataflowPage.tsx b/webui/src/features/dataflow/DataflowPage.tsx index f612115..bc1783d 100644 --- a/webui/src/features/dataflow/DataflowPage.tsx +++ b/webui/src/features/dataflow/DataflowPage.tsx @@ -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([]); - const [draft, setDraft] = useState(null); + const [draft, setDraftValue] = useState(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(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(null); const [selectedNodeId, setSelectedNodeId] = useState(null); const [search, setSearch] = useState(""); @@ -322,15 +347,27 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings }, [savedDraft]); const saveDraft = useCallback(async (): Promise => { + 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) => { - setDraft((current) => current ? { ...current, ...patch } : current); + const current = draftSession.current.value; + if (current) { + const next = { ...current, ...patch }; + draftSession.current.value = next; + setDraftValue(next); + } setSuccess(""); }; @@ -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({ + + +