From c1111c605f88bc7a1bfe66bb71da4ddd0f316ac3 Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Fri, 21 Aug 2026 18:27:09 +0200 Subject: [PATCH] feat(dataflow): govern reusable definition updates --- README.md | 21 +- pyproject.toml | 2 +- src/govoplan_dataflow/__init__.py | 2 +- src/govoplan_dataflow/backend/governance.py | 20 + src/govoplan_dataflow/backend/manifest.py | 16 +- src/govoplan_dataflow/backend/node_library.py | 3 +- src/govoplan_dataflow/backend/router.py | 80 +++- .../backend/schema_validation.py | 52 ++- src/govoplan_dataflow/backend/schemas.py | 21 + src/govoplan_dataflow/backend/service.py | 405 +++++++++++++++++- tests/test_triggers.py | 353 ++++++++++++++- webui/package.json | 2 +- webui/src/api/dataflow.ts | 24 +- webui/src/features/dataflow/DataflowPage.tsx | 233 +++++++++- webui/src/features/dataflow/NodeInspector.tsx | 119 +++-- webui/src/features/dataflow/model.ts | 2 +- webui/src/i18n/generatedTranslations.ts | 38 ++ 17 files changed, 1338 insertions(+), 55 deletions(-) diff --git a/README.md b/README.md index eefa12f..7ed7c23 100644 --- a/README.md +++ b/README.md @@ -137,6 +137,13 @@ records the effective Policy decision and ancestor limits. Inherited definitions remain read-only; lower scopes may narrow, but not broaden, execution, reuse, inheritance, or automation permissions. +Derived definitions report when their source has a newer immutable revision; +the source never mutates the child silently. Adopting an update requires the +reviewed source revision and hash plus a reason. It appends a new child +revision, retains the previous graph and all run evidence, records reviewer +and Policy provenance, and returns the child to draft before the changed graph +can run or receive automation. + Complete active flows support explicit user/API starts, administrative backfills, one-time schedules, interval schedules, and exact-match platform events. Trigger deliveries are durable and idempotent. They enqueue the same @@ -147,11 +154,15 @@ the run before source access or output publication. Confidential and restricted events are not accepted through the direct ingress; those require Core's transactional event bridge. -Reusable subflow nodes pin a template reference, version, graph snapshot, and -parameter values. Their single input is bound to an explicitly marked inline -source inside the snapshot, parameter substitution is data-only, and nesting -is bounded. This keeps completed run definitions reproducible even when the -source template changes later. +Reusable subflow nodes select a Policy-authorized complete flow or template and +an immutable revision. The server resolves the graph instead of accepting a +caller-supplied snapshot, records the source hash and Policy decision, and pins +closed typed input/output contracts. Their single input is bound to an +explicitly marked typed inline source inside the snapshot, parameter +substitution is data-only, and cycles across nested references are rejected. +Incompatible caller schemas fail validation before execution. This keeps +completed run definitions reproducible even when the source definition changes +later. The executable fixtures in `fixtures/golden` cover monthly structured-file reconciliation, sanctions screening, a HEICO-style current-status export, and diff --git a/pyproject.toml b/pyproject.toml index f54d256..3a7786c 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "govoplan-dataflow" -version = "0.1.19" +version = "0.1.20" description = "Governed graphical and SQL data pipelines for GovOPlaN." readme = "README.md" requires-python = ">=3.12" diff --git a/src/govoplan_dataflow/__init__.py b/src/govoplan_dataflow/__init__.py index 260ef6f..3225ed0 100644 --- a/src/govoplan_dataflow/__init__.py +++ b/src/govoplan_dataflow/__init__.py @@ -1,3 +1,3 @@ from __future__ import annotations -__version__ = "0.1.19" +__version__ = "0.1.20" diff --git a/src/govoplan_dataflow/backend/governance.py b/src/govoplan_dataflow/backend/governance.py index 3af940a..d96dcea 100644 --- a/src/govoplan_dataflow/backend/governance.py +++ b/src/govoplan_dataflow/backend/governance.py @@ -121,6 +121,7 @@ def definition_governance_payload( *, principal: ApiPrincipal, registry: object | None, + source_update: Mapping[str, object] | None = None, ) -> dict[str, object]: actions = { action: definition_decision( @@ -142,6 +143,25 @@ def definition_governance_payload( "derived_from_pipeline_id": pipeline.derived_from_pipeline_id, "derived_from_revision": pipeline.derived_from_revision, "derived_from_hash": pipeline.derived_from_hash, + "source_available": bool( + source_update and source_update.get("source_available") is True + ), + "source_name": ( + source_update.get("source_name") if source_update else None + ), + "source_current_revision": ( + source_update.get("source_current_revision") + if source_update + else None + ), + "source_current_hash": ( + source_update.get("source_current_hash") + if source_update + else None + ), + "update_available": bool( + source_update and source_update.get("update_available") is True + ), "derivation_provenance": dict(pipeline.derivation_provenance), "actions": actions, } diff --git a/src/govoplan_dataflow/backend/manifest.py b/src/govoplan_dataflow/backend/manifest.py index 2cf9eef..a9f256d 100644 --- a/src/govoplan_dataflow/backend/manifest.py +++ b/src/govoplan_dataflow/backend/manifest.py @@ -56,7 +56,7 @@ from govoplan_dataflow.backend.dsar_provider import ( MODULE_ID = "dataflow" MODULE_NAME = "Dataflow" -MODULE_VERSION = "0.1.19" +MODULE_VERSION = "0.1.20" READ_SCOPE = "dataflow:pipeline:read" WRITE_SCOPE = "dataflow:pipeline:write" @@ -227,6 +227,9 @@ DOCUMENTATION = ( "Every graph node declares typed inputs, configuration, output schema, and validation rules. " "Source nodes pin inline content or governed Datasource references; combine, filter, transform, " "quality, reconciliation, reusable-subflow, and output nodes remain explicit in the canonical graph. " + "Reusable subflows select a Policy-authorized immutable flow or template revision. The server resolves " + "and pins its graph, source hash, Policy decision, and closed typed input/output contracts; caller-supplied " + "graph snapshots are ignored, incompatible inputs fail validation, and nested reference cycles are rejected. " "Reconciliation rows expose stable key hashes, explicit before/after values, and input hashes. The " "review dialog records accept, reject, correct, or defer decisions in tenant-owned immutable decision " "sets. Their current projection is a fingerprinted Dataflow source; every superseded revision retains " @@ -249,6 +252,9 @@ DOCUMENTATION = ( "dataflow.field.source", "dataflow.field.expression", "dataflow.field.schema", + "dataflow.field.reusable-input-binding", + "dataflow.field.subflow-reference", + "dataflow.field.subflow-revision", "dataflow.action.preview-node", "dataflow.action.review-decisions", ], @@ -261,7 +267,10 @@ DOCUMENTATION = ( body=( "Scope determines ownership and Policy inheritance. Templates can be derived but not run; complete " "flows may be previewed, revisioned, automated, and executed when effective Policy allows it. Saving " - "appends an immutable revision. A scoped copy pins its source revision and content hash. Triggers pin " + "appends an immutable revision. A scoped copy pins its source revision and content hash. A newer " + "source revision is reported without changing the copy. Adopting it requires an exact reviewed source " + "hash and a reason, appends an immutable copy revision, records the Policy decision and reviewer, and " + "returns the copy to draft so runs and automation cannot use the changed graph before activation. Triggers pin " "the revision and authorization grant, then re-evaluate authority for every delivery. Allow runs is " "the definition-level admission boundary and does not grant a caller permission. A one-time run uses " "the configured tenant-local date and time. The missed-run policy either coalesces elapsed interval " @@ -294,8 +303,10 @@ DOCUMENTATION = ( "dataflow.field.trigger-run-at", "dataflow.field.trigger-missed-runs", "dataflow.field.trigger-concurrency", + "dataflow.field.rebase-reason", "dataflow.action.save", "dataflow.action.derive", + "dataflow.action.rebase", "dataflow.action.trigger", "dataflow.action.record-decision", "dataflow.action.delete", @@ -303,6 +314,7 @@ DOCUMENTATION = ( "consequence_classes": { "save_revision": "Appends an immutable pipeline definition revision.", "derive_copy": "Creates a separately governed copy pinned to the source revision and hash.", + "rebase_copy": "Appends the exact reviewed source revision to a scoped copy, records reviewer provenance, and returns it to draft.", "configure_trigger": "Creates or changes an automation command with revision and authorization evidence.", "record_decision": "Appends an actor-attributed immutable decision revision against an exact input hash.", "delete_pipeline": "Prevents future use while retained evidence remains governed.", diff --git a/src/govoplan_dataflow/backend/node_library.py b/src/govoplan_dataflow/backend/node_library.py index ad7e75a..008a7f6 100644 --- a/src/govoplan_dataflow/backend/node_library.py +++ b/src/govoplan_dataflow/backend/node_library.py @@ -613,14 +613,13 @@ _NODE_TYPES = ( type="subflow", category="transform", label="Reusable subflow", - description="Run a pinned parameterized template snapshot as one node.", + description="Run a Policy-authorized, server-resolved immutable definition revision as one node.", icon="boxes", input_ports=(NodePortDefinition(id="input", label="Input"),), config_fields=( NodeConfigField(id="template_ref", label="Template reference", kind="text", required=True), NodeConfigField(id="template_version", label="Template version", kind="text", required=True), NodeConfigField(id="parameters", label="Parameters", kind="json", required=True), - NodeConfigField(id="graph", label="Pinned graph", kind="json", required=True), ), default_config={ "template_ref": "", diff --git a/src/govoplan_dataflow/backend/router.py b/src/govoplan_dataflow/backend/router.py index 2cf5a96..a9e27b0 100644 --- a/src/govoplan_dataflow/backend/router.py +++ b/src/govoplan_dataflow/backend/router.py @@ -68,6 +68,7 @@ from govoplan_dataflow.backend.schemas import ( PipelineRunMetricsResponse, PipelineRunResponse, PipelinePromotionRequest, + PipelineRebaseRequest, PipelineResponse, PipelineSqlResponse, PipelineUpdateRequest, @@ -109,6 +110,7 @@ from govoplan_dataflow.backend.service import ( pipeline_run_response, preview_pipeline, promote_pipeline, + rebase_pipeline, render_graph_sql, start_pipeline_run, update_pipeline, @@ -544,6 +546,8 @@ def api_create_pipeline( tenant_id=tenant_id or principal.tenant_id, actor_id=_actor_id(principal), payload=payload, + principal=principal, + registry=get_registry(), ) except (PermissionError, ValueError) as exc: raise _governance_http_error(exc) from exc @@ -842,6 +846,8 @@ def api_update_pipeline( pipeline_id=pipeline_id, actor_id=_actor_id(principal), payload=payload, + principal=principal, + registry=get_registry(), ) except (PermissionError, ValueError) as exc: raise _governance_http_error(exc) from exc @@ -968,6 +974,64 @@ def api_derive_pipeline( return response +@router.post( + "/pipelines/{pipeline_id}/rebase", + response_model=PipelineResponse, +) +def api_rebase_pipeline( + pipeline_id: str, + payload: PipelineRebaseRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), +) -> PipelineResponse: + _require_any_scope(principal, WRITE_SCOPE, ADMIN_SCOPE) + try: + existing = get_pipeline( + session, + tenant_id=principal.tenant_id, + pipeline_id=pipeline_id, + ) + require_definition_action( + existing, + principal=principal, + registry=get_registry(), + action="edit", + ) + pipeline = rebase_pipeline( + session, + tenant_id=principal.tenant_id, + pipeline_id=pipeline_id, + actor_id=_actor_id(principal), + principal=principal, + registry=get_registry(), + payload=payload, + ) + except PermissionError as exc: + raise _governance_http_error(exc) from exc + except DataflowError as exc: + raise _http_error(exc) from exc + audit_event( + session, + tenant_id=principal.tenant_id, + user_id=getattr(principal.user, "id", None), + api_key_id=principal.api_key_id, + action="dataflow.pipeline.rebased", + object_type="dataflow_pipeline", + object_id=pipeline.id, + details={ + "child_revision": pipeline.current_revision, + "source_pipeline_id": pipeline.derived_from_pipeline_id, + "source_revision": pipeline.derived_from_revision, + "source_hash": pipeline.derived_from_hash, + "status": pipeline.status, + "reason": payload.reason.strip(), + }, + ) + response = _pipeline_response(session, pipeline, principal) + session.commit() + return response + + @router.get( "/pipelines/{pipeline_id}/triggers", response_model=DataflowTriggerListResponse, @@ -1503,10 +1567,22 @@ def api_promote_pipeline( @router.post("/validate", response_model=PipelineValidationResponse) def api_validate_pipeline( payload: PipelineDraftRequest, + session: Session = Depends(get_session), principal: ApiPrincipal = Depends(get_api_principal), ) -> PipelineValidationResponse: _require_any_scope(principal, READ_SCOPE, WRITE_SCOPE, RUN_SCOPE, ADMIN_SCOPE) - return validate_draft(payload) + try: + return validate_draft( + payload, + session=session, + tenant_id=principal.tenant_id, + principal=principal, + registry=get_registry(), + ) + except PermissionError as exc: + raise _governance_http_error(exc) from exc + except DataflowError as exc: + raise _http_error(exc) from exc @router.post("/sql/compile", response_model=PipelineSqlResponse) @@ -1543,6 +1619,8 @@ def api_preview_pipeline( registry=get_registry(), payload=payload, ) + except PermissionError as exc: + raise _governance_http_error(exc) from exc except DataflowError as exc: raise _http_error(exc) from exc if response.pipeline_id: diff --git a/src/govoplan_dataflow/backend/schema_validation.py b/src/govoplan_dataflow/backend/schema_validation.py index dbe5783..236a7b6 100644 --- a/src/govoplan_dataflow/backend/schema_validation.py +++ b/src/govoplan_dataflow/backend/schema_validation.py @@ -129,7 +129,14 @@ def _propagation_context( def _inline_source( context: SchemaPropagationContext, ) -> SchemaPropagationResult: - return SchemaPropagationResult(_inline_schema(context.node.config.get("rows"))) + configured = _configured_schema( + context.node.config.get("contract_schema") + ) + return SchemaPropagationResult( + configured + if configured.columns + else _inline_schema(context.node.config.get("rows")) + ) def _reference_source( @@ -724,16 +731,57 @@ def _comparison_columns(value: object) -> tuple[list[str], list[str]]: def _subflow(context: SchemaPropagationContext) -> SchemaPropagationResult: + input_schema = _configured_schema( + context.node.config.get("input_schema") + ) output_schema = _configured_schema( context.node.config.get("output_schema") ) + diagnostics: list[DataflowDiagnostic] = [] + if input_schema.columns and not context.input_state.open: + missing = sorted(input_schema.columns - context.input_state.columns) + if missing: + diagnostics.append( + _error( + "subflow.input_contract.missing", + "Subflow input is missing required contract columns: " + + ", ".join(missing), + node_id=context.node.id, + field="input_schema", + ) + ) + incompatible = sorted( + column + for column in input_schema.columns & context.input_state.columns + if not _compatible_contract_type( + context.input_state.type_of(column), + input_schema.type_of(column), + ) + ) + if incompatible: + diagnostics.append( + _error( + "subflow.input_contract.type", + "Subflow input has incompatible contract types for: " + + ", ".join(incompatible), + node_id=context.node.id, + field="input_schema", + ) + ) return SchemaPropagationResult( output_schema if output_schema.columns - else unknown_schema() + else unknown_schema(), + tuple(diagnostics), ) +def _compatible_contract_type(actual: str, expected: str) -> bool: + if "unknown" in {actual, expected} or actual == expected: + return True + return {actual, expected} <= {"integer", "number"} + + def _identity(context: SchemaPropagationContext) -> SchemaPropagationResult: return SchemaPropagationResult(context.input_state) diff --git a/src/govoplan_dataflow/backend/schemas.py b/src/govoplan_dataflow/backend/schemas.py index e24b331..1fda064 100644 --- a/src/govoplan_dataflow/backend/schemas.py +++ b/src/govoplan_dataflow/backend/schemas.py @@ -122,6 +122,11 @@ class PipelineGovernanceResponse(BaseModel): derived_from_pipeline_id: str | None derived_from_revision: int | None derived_from_hash: str | None + source_available: bool = False + source_name: str | None = None + source_current_revision: int | None = None + source_current_hash: str | None = None + update_available: bool = False derivation_provenance: dict[str, Any] = Field(default_factory=dict) actions: dict[str, DefinitionActionDecisionResponse] @@ -246,7 +251,23 @@ class PipelineDeriveRequest(BaseModel): allow_automation: bool = False +class PipelineRebaseRequest(BaseModel): + expected_revision: int = Field(ge=1) + source_revision: int = Field(ge=1) + source_hash: str = Field(pattern=r"^[0-9a-f]{64}$") + reason: str = Field(min_length=3, max_length=4_000) + + @field_validator("reason") + @classmethod + def validate_reason(cls, value: str) -> str: + cleaned = value.strip() + if len(cleaned) < 3: + raise ValueError("A meaningful rebase review reason is required.") + return cleaned + + class PipelineDraftRequest(BaseModel): + pipeline_id: str | None = Field(default=None, max_length=36) graph: PipelineGraph | None = None sql_text: str | None = Field(default=None, max_length=100_000) source_nodes: list[GraphNode] = Field(default_factory=list, max_length=20) diff --git a/src/govoplan_dataflow/backend/service.py b/src/govoplan_dataflow/backend/service.py index a344277..fbb525c 100644 --- a/src/govoplan_dataflow/backend/service.py +++ b/src/govoplan_dataflow/backend/service.py @@ -60,6 +60,7 @@ from govoplan_dataflow.backend.graph import ( preserve_compatible_graph_layout, validate_graph, ) +from govoplan_dataflow.backend.ir import graph_to_ir from govoplan_dataflow.backend.schemas import ( DataflowDiagnostic, GraphNode, @@ -72,6 +73,7 @@ from govoplan_dataflow.backend.schemas import ( PipelinePreviewResponse, PipelineDeploymentResponse, PipelinePromotionRequest, + PipelineRebaseRequest, PipelineResponse, PipelineRevisionResponse, PipelineRunResponse, @@ -187,15 +189,177 @@ def get_pipeline_revision( return item +def _resolve_reusable_subflows( + session: Session, + *, + tenant_id: str, + graph: PipelineGraph, + principal: ApiPrincipal | None, + registry: object | None, + target_pipeline_id: str | None, + ancestry: tuple[str, ...] = (), +) -> PipelineGraph: + if not any(node.type == "subflow" for node in graph.nodes): + return graph + if principal is None: + raise DataflowConflictError( + "Reusable subflows require a tenant principal and current Policy " + "decision." + ) + resolved_nodes: list[GraphNode] = [] + for node in graph.nodes: + if node.type != "subflow": + resolved_nodes.append(node) + continue + source_id = _pipeline_id_from_ref(node.config.get("template_ref")) + if source_id == target_pipeline_id: + raise DataflowConflictError( + "A pipeline cannot reference itself as a reusable subflow." + ) + source = get_pipeline( + session, + tenant_id=tenant_id, + pipeline_id=source_id, + ) + reuse_decision = require_definition_action( + source, + principal=principal, + registry=registry, + action="reuse", + ) + source_revision_number = _subflow_revision( + node.config.get("template_version") + ) + source_revision = get_pipeline_revision( + session, + pipeline=source, + revision=source_revision_number, + ) + reference_key = f"{source.id}:{source_revision.revision}" + if reference_key in ancestry: + raise DataflowConflictError( + "Reusable subflow references contain a cycle at " + f"pipeline:{source.id} revision {source_revision.revision}." + ) + nested = _resolve_reusable_subflows( + session, + tenant_id=tenant_id, + graph=PipelineGraph.model_validate(source_revision.graph), + principal=principal, + registry=registry, + target_pipeline_id=target_pipeline_id, + ancestry=(*ancestry, reference_key), + ) + input_nodes = [ + item + for item in nested.nodes + if item.type == "source.inline" + and item.config.get("input_binding") is True + ] + if len(input_nodes) != 1: + raise DataflowConflictError( + "A referenced reusable definition must declare exactly one " + "inline template input binding." + ) + typed = graph_to_ir(nested) + typed_by_id = {item.id: item for item in typed.nodes} + output_nodes = [item for item in nested.nodes if item.type == "output"] + if len(output_nodes) != 1: + raise DataflowConflictError( + "A referenced reusable definition must have exactly one " + "typed output." + ) + input_contract = _typed_contract( + typed_by_id[input_nodes[0].id].output_schema, + label="input", + ) + output_contract = _typed_contract( + typed_by_id[output_nodes[0].id].output_schema, + label="output", + ) + config = { + **node.config, + "template_ref": f"pipeline:{source.id}", + "template_version": str(source_revision.revision), + "template_hash": source_revision.content_hash, + "graph": canonical_graph_payload(nested), + "input_schema": input_contract, + "output_schema": output_contract, + "reference_provenance": { + "source_scope": { + "scope_type": source.scope_type, + "scope_id": source.scope_id, + }, + "source_definition_kind": source.definition_kind, + "policy_decision": reuse_decision.to_dict(), + }, + } + resolved_nodes.append( + node.model_copy(update={"config": config}, deep=True) + ) + return graph.model_copy(update={"nodes": resolved_nodes}, deep=True) + + +def _pipeline_id_from_ref(value: object) -> str: + text = str(value or "").strip() + if not text.startswith("pipeline:") or len(text) <= len("pipeline:"): + raise DataflowConflictError( + "Reusable subflows require a canonical pipeline reference." + ) + return text.removeprefix("pipeline:") + + +def _subflow_revision(value: object) -> int: + try: + revision = int(str(value).strip()) + except (TypeError, ValueError) as exc: + raise DataflowConflictError( + "Reusable subflows require a valid immutable source revision." + ) from exc + if revision < 1: + raise DataflowConflictError( + "Reusable subflow revisions must be positive." + ) + return revision + + +def _typed_contract(schema: object, *, label: str) -> list[dict[str, object]]: + fields = tuple(getattr(schema, "fields", ())) + if not fields or any(getattr(item, "type", "unknown") == "unknown" for item in fields): + raise DataflowConflictError( + f"The reusable definition needs a closed typed {label} contract. " + "Provide representative typed rows at its template input." + ) + return [ + { + "name": str(item.name), + "type": str(item.type), + "nullable": bool(item.nullable), + } + for item in fields + ] + + def create_pipeline( session: Session, *, tenant_id: str, actor_id: str | None, payload: PipelineCreateRequest, + principal: ApiPrincipal | None = None, + registry: object | None = None, ) -> DataflowPipeline: - definition = normalize_definition( + pipeline_id = new_uuid() + graph = _resolve_reusable_subflows( + session, + tenant_id=tenant_id, graph=payload.graph, + principal=principal, + registry=registry, + target_pipeline_id=pipeline_id, + ) + definition = normalize_definition( + graph=graph, sql_text=payload.sql_text, editor_mode=payload.editor_mode, ) @@ -209,6 +373,7 @@ def create_pipeline( else payload.scope_id ) pipeline = DataflowPipeline( + id=pipeline_id, tenant_id=stored_tenant_id, scope_type=payload.scope_type, scope_id=scope_id, @@ -248,6 +413,8 @@ def update_pipeline( pipeline_id: str, actor_id: str | None, payload: PipelineUpdateRequest, + principal: ApiPrincipal | None = None, + registry: object | None = None, ) -> DataflowPipeline: pipeline = get_pipeline(session, tenant_id=tenant_id, pipeline_id=pipeline_id) if payload.expected_revision != pipeline.current_revision: @@ -270,8 +437,16 @@ def update_pipeline( raise DataflowConflictError( "Definition kind is immutable; derive a flow or template instead." ) - definition = normalize_definition( + graph = _resolve_reusable_subflows( + session, + tenant_id=tenant_id, graph=payload.graph, + principal=principal, + registry=registry, + target_pipeline_id=pipeline.id, + ) + definition = normalize_definition( + graph=graph, sql_text=payload.sql_text, editor_mode=payload.editor_mode, ) @@ -417,6 +592,195 @@ def derive_pipeline( return pipeline +def pipeline_source_update_status( + session: Session, + *, + tenant_id: str, + pipeline: DataflowPipeline, +) -> dict[str, object]: + source_id = pipeline.derived_from_pipeline_id + if not source_id: + return { + "source_available": False, + "source_name": None, + "source_current_revision": None, + "source_current_hash": None, + "update_available": False, + } + source = session.scalar( + select(DataflowPipeline).where( + DataflowPipeline.id == source_id, + or_( + DataflowPipeline.tenant_id == tenant_id, + DataflowPipeline.tenant_id.is_(None), + ), + DataflowPipeline.deleted_at.is_(None), + ) + ) + if source is None: + return { + "source_available": False, + "source_name": None, + "source_current_revision": None, + "source_current_hash": None, + "update_available": False, + } + revision = get_pipeline_revision(session, pipeline=source) + return { + "source_available": True, + "source_name": source.name, + "source_current_revision": revision.revision, + "source_current_hash": revision.content_hash, + "update_available": ( + revision.revision != pipeline.derived_from_revision + or revision.content_hash != pipeline.derived_from_hash + ), + } + + +def rebase_pipeline( + session: Session, + *, + tenant_id: str, + pipeline_id: str, + actor_id: str | None, + principal: ApiPrincipal, + registry: object | None, + payload: PipelineRebaseRequest, +) -> DataflowPipeline: + pipeline = get_pipeline( + session, + tenant_id=tenant_id, + pipeline_id=pipeline_id, + ) + if payload.expected_revision != pipeline.current_revision: + raise DataflowConflictError( + "Derived pipeline changed on the server; expected revision " + f"{payload.expected_revision}, current revision is " + f"{pipeline.current_revision}." + ) + source_id = pipeline.derived_from_pipeline_id + if not source_id: + raise DataflowConflictError( + "Only a pipeline derived from another definition can adopt a " + "source update." + ) + source = get_pipeline( + session, + tenant_id=tenant_id, + pipeline_id=source_id, + ) + reuse_decision = require_definition_action( + source, + principal=principal, + registry=registry, + action="derive", + ) + source_revision = get_pipeline_revision( + session, + pipeline=source, + revision=payload.source_revision, + ) + if source_revision.content_hash != payload.source_hash: + raise DataflowConflictError( + "The reviewed source hash no longer matches the requested " + "revision; reload before adopting the update." + ) + if ( + pipeline.derived_from_revision is not None + and source_revision.revision <= pipeline.derived_from_revision + ): + raise DataflowConflictError( + "A source update must use a revision newer than the currently " + "pinned revision." + ) + + previous_child_revision = pipeline.current_revision + previous_child = get_pipeline_revision(session, pipeline=pipeline) + previous_source_revision = pipeline.derived_from_revision + previous_source_hash = pipeline.derived_from_hash + source_limits = _effective_governance_limits( + source, + decision_details=reuse_decision.details, + ) + effective_limits = { + "inherit_to_lower_scopes": ( + pipeline.inherit_to_lower_scopes + and source_limits["inherit_to_lower_scopes"] + ), + "allow_run": pipeline.allow_run and source_limits["allow_run"], + "allow_reuse": pipeline.allow_reuse and source_limits["allow_reuse"], + "allow_automation": ( + pipeline.allow_automation and source_limits["allow_automation"] + ), + } + next_child_revision = previous_child_revision + 1 + rebased_at = utcnow() + history_value = pipeline.derivation_provenance.get("rebase_history", []) + history = list(history_value) if isinstance(history_value, list) else [] + history.append( + { + "child_revision_before": previous_child_revision, + "child_hash_before": previous_child.content_hash, + "child_revision_after": next_child_revision, + "source_revision_before": previous_source_revision, + "source_hash_before": previous_source_hash, + "source_revision_after": source_revision.revision, + "source_hash_after": source_revision.content_hash, + "policy_decision": reuse_decision.to_dict(), + "reason": payload.reason.strip(), + "rebased_by": actor_id, + "rebased_at": rebased_at.isoformat(), + } + ) + provenance = dict(pipeline.derivation_provenance) + provenance.update( + { + "source_ref": f"pipeline:{source.id}", + "source_scope": { + "scope_type": source.scope_type, + "scope_id": source.scope_id, + }, + "source_definition_kind": source.definition_kind, + "source_revision": source_revision.revision, + "source_hash": source_revision.content_hash, + "source_effective_limits": effective_limits, + "policy_decision": reuse_decision.to_dict(), + "last_rebased_by": actor_id, + "last_rebased_at": rebased_at.isoformat(), + "last_rebase_reason": payload.reason.strip(), + "rebase_history": history, + } + ) + + pipeline.current_revision = next_child_revision + pipeline.status = "draft" + pipeline.inherit_to_lower_scopes = effective_limits[ + "inherit_to_lower_scopes" + ] + pipeline.allow_run = effective_limits["allow_run"] + pipeline.allow_reuse = effective_limits["allow_reuse"] + pipeline.allow_automation = effective_limits["allow_automation"] + pipeline.derived_from_revision = source_revision.revision + pipeline.derived_from_hash = source_revision.content_hash + pipeline.derivation_provenance = provenance + pipeline.updated_by = actor_id + pipeline.revisions.append( + DataflowPipelineRevision( + tenant_id=pipeline.tenant_id, + revision=next_child_revision, + schema_version=source_revision.schema_version, + graph=json.loads(json.dumps(source_revision.graph)), + sql_text=source_revision.sql_text, + editor_mode=source_revision.editor_mode, + content_hash=source_revision.content_hash, + created_by=actor_id, + ) + ) + session.flush() + return pipeline + + def delete_pipeline( session: Session, *, @@ -455,11 +819,36 @@ def pipeline_response( pipeline, principal=principal, registry=registry, + source_update=pipeline_source_update_status( + session, + tenant_id=principal.tenant_id, + pipeline=pipeline, + ), ), ) -def validate_draft(payload: PipelineDraftRequest) -> PipelineValidationResponse: +def validate_draft( + payload: PipelineDraftRequest, + *, + session: Session | None = None, + tenant_id: str | None = None, + principal: ApiPrincipal | None = None, + registry: object | None = None, +) -> PipelineValidationResponse: + if payload.graph is not None and session is not None and tenant_id is not None: + payload = payload.model_copy( + update={ + "graph": _resolve_reusable_subflows( + session, + tenant_id=tenant_id, + graph=payload.graph, + principal=principal, + registry=registry, + target_pipeline_id=payload.pipeline_id, + ) + } + ) if payload.sql_text and payload.sql_text.strip(): try: graph, sql_text, diagnostics = compile_sql( @@ -598,7 +987,13 @@ def preview_pipeline( sql_text=payload.sql_text, source_nodes=payload.source_nodes, ) - validated = validate_draft(draft) + validated = validate_draft( + draft, + session=session, + tenant_id=tenant_id, + principal=principal, + registry=registry, + ) if not validated.valid or validated.graph is None: return PipelinePreviewResponse( run_id=None, @@ -2290,11 +2685,13 @@ __all__ = [ "list_pipelines", "normalize_definition", "pipeline_response", + "pipeline_source_update_status", "pipeline_deployment_response", "pipeline_run_descriptor", "pipeline_run_request", "pipeline_run_response", "promote_pipeline", + "rebase_pipeline", "preview_pipeline", "render_graph_sql", "start_pipeline_run", diff --git a/tests/test_triggers.py b/tests/test_triggers.py index bbd5da0..fe11a9c 100644 --- a/tests/test_triggers.py +++ b/tests/test_triggers.py @@ -25,12 +25,16 @@ from govoplan_dataflow.backend.schemas import ( DataflowTriggerSchedule, PipelineCreateRequest, PipelineDeriveRequest, + PipelineRebaseRequest, PipelineUpdateRequest, ) from govoplan_dataflow.backend.service import ( DataflowConflictError, + DataflowValidationError, create_pipeline, derive_pipeline, + pipeline_response, + rebase_pipeline, start_pipeline_run, update_pipeline, ) @@ -46,10 +50,87 @@ POLICY_CAPABILITY = "policy.definitionGovernance" AUTOMATION_CAPABILITY = "auth.automationPrincipalProvider" -def sample_graph(): +def sample_graph(*, minimum: int = 10): from test_service import sample_graph as build_graph - return build_graph() + return build_graph(minimum=minimum) + + +def reusable_graph(*, minimum: int = 10): + graph = sample_graph(minimum=minimum) + return graph.model_copy( + update={ + "nodes": [ + ( + node.model_copy( + update={ + "config": { + **node.config, + "input_binding": True, + } + }, + deep=True, + ) + if node.id == "source" + else node + ) + for node in graph.nodes + ] + }, + deep=True, + ) + + +def referencing_graph( + source_ref: str, + source_revision: int, + *, + omit_amount: bool = False, + input_binding: bool = False, +): + graph = sample_graph() + return graph.model_copy( + update={ + "nodes": [ + ( + node.model_copy( + update={ + "config": { + **node.config, + "rows": ( + [{"id": 1}] + if omit_amount + else node.config["rows"] + ), + "input_binding": input_binding, + } + }, + deep=True, + ) + if node.id == "source" + else node.model_copy( + update={ + "type": "subflow", + "label": "Governed reusable flow", + "config": { + "template_ref": source_ref, + "template_version": str(source_revision), + "parameters": {}, + "graph": sample_graph( + minimum=999 + ).model_dump(mode="json"), + }, + }, + deep=True, + ) + if node.id == "filter" + else node + ) + for node in graph.nodes + ] + }, + deep=True, + ) def principal() -> ApiPrincipal: @@ -482,6 +563,152 @@ class DataflowTriggerTests(unittest.TestCase): ), ) + def test_reusable_reference_is_policy_resolved_with_typed_contracts( + self, + ) -> None: + template = create_pipeline( + self.session, + tenant_id="tenant-1", + actor_id="template-author", + payload=PipelineCreateRequest( + name="Typed reusable filter", + graph=reusable_graph(minimum=10), + definition_kind="template", + allow_reuse=True, + ), + ) + consumer = create_pipeline( + self.session, + tenant_id="tenant-1", + actor_id="account-1", + principal=self.actor, + registry=self.registry, + payload=PipelineCreateRequest( + name="Resolved consumer", + graph=referencing_graph(f"pipeline:{template.id}", 1), + ), + ) + self.session.commit() + + stored = consumer.revisions[0].graph + subflow = next( + node for node in stored["nodes"] if node["id"] == "filter" + ) + self.assertEqual(template.revisions[0].content_hash, subflow["config"]["template_hash"]) + self.assertEqual( + 10, + next( + node + for node in subflow["config"]["graph"]["nodes"] + if node["id"] == "filter" + )["config"]["value"], + ) + self.assertEqual( + {"id", "amount"}, + { + field["name"] + for field in subflow["config"]["input_schema"] + }, + ) + self.assertEqual( + {"id", "amount"}, + { + field["name"] + for field in subflow["config"]["output_schema"] + }, + ) + self.assertTrue( + subflow["config"]["reference_provenance"][ + "policy_decision" + ]["allowed"] + ) + + with self.assertRaises(DataflowValidationError): + create_pipeline( + self.session, + tenant_id="tenant-1", + actor_id="account-1", + principal=self.actor, + registry=self.registry, + payload=PipelineCreateRequest( + name="Incompatible consumer", + graph=referencing_graph( + f"pipeline:{template.id}", + 1, + omit_amount=True, + ), + ), + ) + + def test_reusable_reference_cycles_are_rejected_across_revisions( + self, + ) -> None: + left = create_pipeline( + self.session, + tenant_id="tenant-1", + actor_id="account-1", + payload=PipelineCreateRequest( + name="Left template", + graph=reusable_graph(), + definition_kind="template", + allow_reuse=True, + ), + ) + right = create_pipeline( + self.session, + tenant_id="tenant-1", + actor_id="account-1", + payload=PipelineCreateRequest( + name="Right template", + graph=reusable_graph(), + definition_kind="template", + allow_reuse=True, + ), + ) + self.session.flush() + update_pipeline( + self.session, + tenant_id="tenant-1", + pipeline_id=left.id, + actor_id="account-1", + principal=self.actor, + registry=self.registry, + payload=PipelineUpdateRequest( + name=left.name, + graph=referencing_graph( + f"pipeline:{right.id}", + 1, + input_binding=True, + ), + status="draft", + expected_revision=1, + definition_kind="template", + allow_reuse=True, + ), + ) + + with self.assertRaisesRegex(DataflowConflictError, "cannot reference itself"): + update_pipeline( + self.session, + tenant_id="tenant-1", + pipeline_id=right.id, + actor_id="account-1", + principal=self.actor, + registry=self.registry, + payload=PipelineUpdateRequest( + name=right.name, + graph=referencing_graph( + f"pipeline:{left.id}", + 2, + input_binding=True, + ), + status="draft", + expected_revision=1, + definition_kind="template", + allow_reuse=True, + ), + ) + def test_derived_limits_cannot_be_broadened_transitively(self) -> None: template = create_pipeline( self.session, @@ -549,6 +776,128 @@ class DataflowTriggerTests(unittest.TestCase): self.assertFalse(grandchild.allow_automation) self.assertFalse(grandchild.inherit_to_lower_scopes) + def test_source_update_is_detected_and_rebased_as_reviewed_revision( + self, + ) -> None: + template = create_pipeline( + self.session, + tenant_id="tenant-1", + actor_id="account-1", + payload=PipelineCreateRequest( + name="Reusable import", + graph=sample_graph(), + definition_kind="template", + allow_reuse=True, + allow_automation=True, + ), + ) + derived = derive_pipeline( + self.session, + tenant_id="tenant-1", + actor_id="account-1", + principal=self.actor, + registry=self.registry, + source_pipeline_id=template.id, + payload=PipelineDeriveRequest( + name="Tenant import", + allow_run=True, + allow_automation=True, + ), + ) + self.session.commit() + + before = pipeline_response( + self.session, + derived, + principal=self.actor, + registry=self.registry, + ) + self.assertTrue(before.governance.source_available) + self.assertFalse(before.governance.update_available) + original_child_hash = derived.revisions[0].content_hash + + update_pipeline( + self.session, + tenant_id="tenant-1", + pipeline_id=template.id, + actor_id="template-author", + payload=PipelineUpdateRequest( + name=template.name, + graph=sample_graph(minimum=20), + status="draft", + expected_revision=1, + definition_kind="template", + allow_reuse=True, + allow_automation=True, + ), + ) + self.session.commit() + source_hash = template.revisions[-1].content_hash + + available = pipeline_response( + self.session, + derived, + principal=self.actor, + registry=self.registry, + ) + self.assertTrue(available.governance.update_available) + self.assertEqual(2, available.governance.source_current_revision) + self.assertEqual(source_hash, available.governance.source_current_hash) + self.assertEqual(original_child_hash, derived.revisions[0].content_hash) + + with self.assertRaisesRegex(DataflowConflictError, "source hash"): + rebase_pipeline( + self.session, + tenant_id="tenant-1", + pipeline_id=derived.id, + actor_id="reviewer-1", + principal=self.actor, + registry=self.registry, + payload=PipelineRebaseRequest( + expected_revision=1, + source_revision=2, + source_hash="0" * 64, + reason="Reviewed the changed filter threshold.", + ), + ) + + rebased = rebase_pipeline( + self.session, + tenant_id="tenant-1", + pipeline_id=derived.id, + actor_id="reviewer-1", + principal=self.actor, + registry=self.registry, + payload=PipelineRebaseRequest( + expected_revision=1, + source_revision=2, + source_hash=source_hash, + reason="Reviewed the changed filter threshold.", + ), + ) + self.session.commit() + + self.assertEqual(2, rebased.current_revision) + self.assertEqual("draft", rebased.status) + self.assertEqual(2, rebased.derived_from_revision) + self.assertEqual(source_hash, rebased.derived_from_hash) + self.assertEqual(source_hash, rebased.revisions[-1].content_hash) + self.assertEqual(original_child_hash, rebased.revisions[0].content_hash) + history = rebased.derivation_provenance["rebase_history"] + self.assertEqual(1, len(history)) + self.assertEqual("reviewer-1", history[0]["rebased_by"]) + self.assertEqual( + "Reviewed the changed filter threshold.", + history[0]["reason"], + ) + current = pipeline_response( + self.session, + rebased, + principal=self.actor, + registry=self.registry, + ) + self.assertFalse(current.governance.update_available) + if __name__ == "__main__": unittest.main() diff --git a/webui/package.json b/webui/package.json index 0a3b511..6b973dd 100644 --- a/webui/package.json +++ b/webui/package.json @@ -1,6 +1,6 @@ { "name": "@govoplan/dataflow-webui", - "version": "0.1.19", + "version": "0.1.20", "private": true, "type": "module", "main": "src/index.ts", diff --git a/webui/src/api/dataflow.ts b/webui/src/api/dataflow.ts index 11adba1..8c37377 100644 --- a/webui/src/api/dataflow.ts +++ b/webui/src/api/dataflow.ts @@ -95,6 +95,11 @@ export type PipelineGovernance = { derived_from_pipeline_id?: string | null; derived_from_revision?: number | null; derived_from_hash?: string | null; + source_available: boolean; + source_name?: string | null; + source_current_revision?: number | null; + source_current_hash?: string | null; + update_available: boolean; derivation_provenance: Record; actions: Record; }; @@ -505,6 +510,23 @@ export function deriveDataflowPipeline( ); } +export function rebaseDataflowPipeline( + settings: ApiSettings, + pipelineId: string, + payload: { + expected_revision: number; + source_revision: number; + source_hash: string; + reason: string; + } +): Promise { + return apiFetch( + settings, + `/api/v1/dataflow/pipelines/${encodeURIComponent(pipelineId)}/rebase`, + { method: "POST", body: JSON.stringify(payload) } + ); +} + export function dataflowScopeReferenceProvider( settings: ApiSettings, scopeType: "user" | "group" @@ -564,7 +586,7 @@ export function deleteDataflowTrigger( export function validateDataflowPipeline( settings: ApiSettings, - payload: { graph?: PipelineGraph; sql_text?: string; source_nodes?: PipelineGraphNode[] } + payload: { pipeline_id?: string | null; graph?: PipelineGraph; sql_text?: string; source_nodes?: PipelineGraphNode[] } ): Promise { return apiFetch(settings, "/api/v1/dataflow/validate", { method: "POST", diff --git a/webui/src/features/dataflow/DataflowPage.tsx b/webui/src/features/dataflow/DataflowPage.tsx index 0d82129..3f1bda4 100644 --- a/webui/src/features/dataflow/DataflowPage.tsx +++ b/webui/src/features/dataflow/DataflowPage.tsx @@ -13,6 +13,7 @@ import { Code2, CopyPlus, DatabaseZap, + GitCompareArrows, ListChecks, Network, Play, @@ -84,6 +85,7 @@ import { listDataflowTriggers, previewDataflowPipeline, promoteDataflowPipeline, + rebaseDataflowPipeline, recordDataflowDecision, runDataflowPipeline, renderDataflowSql, @@ -161,6 +163,7 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings const [runOpen, setRunOpen] = useState(false); const [definitionSettingsOpen, setDefinitionSettingsOpen] = useState(false); const [deriveOpen, setDeriveOpen] = useState(false); + const [rebaseOpen, setRebaseOpen] = useState(false); const [triggersOpen, setTriggersOpen] = useState(false); const [decisionReviewOpen, setDecisionReviewOpen] = useState(false); const [nodeLibrary, setNodeLibrary] = useState(FALLBACK_NODE_LIBRARY); @@ -183,6 +186,14 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings && canWrite && draft.governance?.actions.derive?.allowed ); + const canRebase = Boolean( + draft?.id + && draft.governance?.update_available + && draft.governance.source_current_revision + && draft.governance.source_current_hash + && canEdit + && !dirty + ); const canStartSavedRun = Boolean( draft?.id && draft.definitionKind === "flow" @@ -411,10 +422,15 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings setSuccess(""); try { const response = await validateDataflowPipeline(settings, draft.editorMode === "sql" - ? { graph: draft.graph, sql_text: draft.sqlText, source_nodes: sourceNodes(draft.graph) } - : { graph: draft.graph }); + ? { pipeline_id: draft.id, graph: draft.graph, sql_text: draft.sqlText, source_nodes: sourceNodes(draft.graph) } + : { pipeline_id: draft.id, graph: draft.graph }); setDiagnostics(response.diagnostics); - if (response.valid) setSuccess("Pipeline definition is valid."); + if (response.valid) { + if (response.graph && draft.editorMode === "graph") { + updateDraft({ graph: response.graph }); + } + setSuccess("Pipeline definition is valid."); + } setResultOpen(true); setResultTab("diagnostics"); } catch (validationError) { @@ -766,6 +782,26 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings } /> ) : null} + {draft.governance?.derived_from_pipeline_id ? ( + } + variant="ghost" + onClick={() => setRebaseOpen(true)} + disabled={!canRebase} + disabledReason={ + dirty + ? DATAFLOW_I18N.saveFirst + : !canEdit + ? editBlockedReason + : !draft.governance.source_available + ? "The source definition is no longer available." + : !draft.governance.update_available + ? "This copy already pins the current source revision." + : undefined + } + /> + ) : null} {draft.id ? ( ( + pipeline.id !== draft.id + && pipeline.governance.actions.reuse?.allowed + ))} sources={sources} sourceCatalogueAvailable={sourceCatalogueAvailable} readOnly={!canEdit} @@ -1060,6 +1100,32 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings setSuccess("Created a pinned scoped copy."); }} /> + setRebaseOpen(false)} + onRebased={(pipeline) => { + const next = draftFromPipeline(pipeline); + setPipelines((current) => [ + pipeline, + ...current.filter((item) => item.id !== pipeline.id) + ]); + setDraft(next); + setSavedDraft(structuredClone(next)); + setSelectedNodeId(next.graph.nodes[0]?.id ?? null); + setRebaseOpen(false); + setPreview(null); + setDiagnostics([]); + setNodeDiagnostics([]); + setSuccess(`Adopted source revision ${pipeline.governance.derived_from_revision} as draft revision ${pipeline.current_revision}.`); + }} + /> {draft.governance.derived_from_hash} + {!draft.governance.source_available ? ( + + ) : draft.governance.update_available ? ( + <> + + + {draft.governance.source_name ?? "Source definition"} + {" · revision "} + {draft.governance.source_current_revision} + + + ) : ( + + )} ) : null} {provenance.length ? ( @@ -1430,6 +1510,153 @@ function DerivePipelineDialog({ ); } +function RebasePipelineDialog({ + open, + settings, + pipeline, + onClose, + onRebased +}: { + open: boolean; + settings: ApiSettings; + pipeline: { + id: string; + name: string; + currentRevision: number; + governance: Pipeline["governance"]; + } | null; + onClose: () => void; + onRebased: (pipeline: Pipeline) => void; +}) { + const { requestDiscard } = useUnsavedChanges(); + const [reason, setReason] = useState(""); + const [busy, setBusy] = useState(false); + const [error, setError] = useState(""); + const governance = pipeline?.governance; + const sourceRevision = governance?.source_current_revision ?? null; + const sourceHash = governance?.source_current_hash ?? null; + const dirty = Boolean(open && reason); + + useEffect(() => { + if (!open) return; + setReason(""); + setError(""); + }, [open, pipeline?.id, sourceRevision]); + + const resetDraft = () => { + setReason(""); + setError(""); + }; + + const rebase = async (): Promise => { + if ( + !pipeline + || !sourceRevision + || !sourceHash + || !reason.trim() + || !governance?.update_available + ) return false; + setBusy(true); + setError(""); + try { + onRebased(await rebaseDataflowPipeline(settings, pipeline.id, { + expected_revision: pipeline.currentRevision, + source_revision: sourceRevision, + source_hash: sourceHash, + reason: reason.trim() + })); + return true; + } catch (rebaseError) { + setError(apiErrorMessage(rebaseError)); + return false; + } finally { + setBusy(false); + } + }; + + useUnsavedDraftGuard({ + dirty, + title: "Unapplied source update", + message: "Apply the reviewed source update or discard the review reason before leaving.", + onSave: rebase, + onDiscard: resetDraft + }); + + const close = () => { + if (busy) return; + if (dirty) requestDiscard(onClose); + else onClose(); + }; + + return ( + + + + + )} + > +
+ {error ? {error} : null} + + {governance?.source_name ?? "Source definition"} + + Pinned revision {governance?.derived_from_revision ?? "—"} + {" → source revision "} + {sourceRevision ?? "—"} + + {sourceHash ? {sourceHash} : null} + + + Adopting the update replaces the copy's current graph with the exact + reviewed source revision and returns the copy to draft. Existing + revisions, run evidence and rebase provenance remain immutable. + + +