diff --git a/src/govoplan_dataflow/backend/governance.py b/src/govoplan_dataflow/backend/governance.py index 09cabef..3af940a 100644 --- a/src/govoplan_dataflow/backend/governance.py +++ b/src/govoplan_dataflow/backend/governance.py @@ -67,7 +67,7 @@ def normalize_definition_scope( return principal.tenant_id, "group", clean_id if clean_type == "user": if not clean_id: - clean_id = principal.membership_id or principal.account_id + clean_id = principal.account_id if ( clean_id not in {principal.membership_id, principal.account_id} and not administrative @@ -75,6 +75,8 @@ def normalize_definition_scope( raise PermissionError( "Definitions can only be created for the current user." ) + if clean_id == principal.membership_id: + clean_id = principal.account_id return principal.tenant_id, "user", clean_id raise ValueError( "Definition scope must be system, tenant, group, or user." diff --git a/src/govoplan_dataflow/backend/manifest.py b/src/govoplan_dataflow/backend/manifest.py index d6b2c83..1c90789 100644 --- a/src/govoplan_dataflow/backend/manifest.py +++ b/src/govoplan_dataflow/backend/manifest.py @@ -7,6 +7,7 @@ from govoplan_core.core.module_guards import ( persistent_table_uninstall_guard, ) from govoplan_core.core.access import ( + CAPABILITY_ACCESS_DIRECTORY, CAPABILITY_AUTH_AUTOMATION_PRINCIPAL_PROVIDER, ) from govoplan_core.core.dataflows import ( @@ -34,6 +35,7 @@ from govoplan_core.core.datasources import ( from govoplan_core.core.policy import ( CAPABILITY_POLICY_DEFINITION_GOVERNANCE, ) +from govoplan_core.core.views import ViewSurface from govoplan_core.db.base import Base from govoplan_dataflow.backend.db import models as dataflow_models @@ -241,6 +243,7 @@ manifest = ModuleManifest( "workflow", ), optional_capabilities=( + CAPABILITY_ACCESS_DIRECTORY, CAPABILITY_DATASOURCE_CATALOGUE, CAPABILITY_DATASOURCE_LIFECYCLE, CAPABILITY_DATASOURCE_PUBLICATION, @@ -320,6 +323,15 @@ manifest = ModuleManifest( order=72, ), ), + view_surfaces=( + ViewSurface( + id="dataflow.widget.pipelines", + module_id=MODULE_ID, + kind="section", + label="Dataflows widget", + order=70, + ), + ), ), route_factory=_dataflow_router, capability_factories={ diff --git a/src/govoplan_dataflow/backend/router.py b/src/govoplan_dataflow/backend/router.py index ac37514..d2bd953 100644 --- a/src/govoplan_dataflow/backend/router.py +++ b/src/govoplan_dataflow/backend/router.py @@ -1,8 +1,12 @@ from __future__ import annotations -from fastapi import APIRouter, Depends, HTTPException, status +from fastapi import APIRouter, Depends, HTTPException, Query, status from sqlalchemy.orm import Session +from govoplan_core.api.v1.schemas import ( + ReferenceOptionListResponse, + ReferenceOptionResponse, +) from govoplan_core.audit.logging import audit_event from govoplan_core.auth import ApiPrincipal, get_api_principal, has_scope from govoplan_core.core.automation import AutomationInvocation @@ -21,6 +25,11 @@ from govoplan_core.core.datasources import ( datasource_catalogue, datasource_lifecycle, ) +from govoplan_core.core.references import ( + access_scope_reference_options, + access_scope_reference_provider_available, + validate_access_scope_reference, +) from govoplan_core.core.tabular_sources import ( parse_tabular_csv, ) @@ -264,6 +273,37 @@ def api_node_types( return _node_library_response() +@router.get("/scope-targets", response_model=ReferenceOptionListResponse) +def api_scope_targets( + scope_type: str, + q: str = "", + selected: list[str] = Query(default=[]), + limit: int = Query(default=50, ge=1, le=200), + principal: ApiPrincipal = Depends(get_api_principal), +) -> ReferenceOptionListResponse: + _require_any_scope(principal, READ_SCOPE, WRITE_SCOPE, ADMIN_SCOPE) + registry = get_registry() + try: + options = access_scope_reference_options( + registry, + principal, + scope_type=scope_type, + query=q, + selected_values=selected, + limit=limit, + administrative=has_scope(principal, ADMIN_SCOPE), + ) + except ValueError as exc: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, + detail=str(exc), + ) from exc + return ReferenceOptionListResponse( + options=[ReferenceOptionResponse(**item.to_dict()) for item in options], + provider_available=access_scope_reference_provider_available(registry), + ) + + @router.get("/sources", response_model=TabularSourceListResponse) def api_list_sources( query: str = "", @@ -417,6 +457,13 @@ def api_create_pipeline( payload = payload.model_copy( update={"scope_type": scope_type, "scope_id": scope_id} ) + scope_id = validate_access_scope_reference( + get_registry(), + tenant_id=tenant_id or principal.tenant_id, + scope_type=scope_type, + scope_id=scope_id, + ) + payload = payload.model_copy(update={"scope_id": scope_id}) pipeline = create_pipeline( session, tenant_id=tenant_id or principal.tenant_id, @@ -488,6 +535,26 @@ def api_update_pipeline( registry=get_registry(), action="edit", ) + tenant_id, scope_type, scope_id = normalize_definition_scope( + principal, + scope_type=payload.scope_type, + scope_id=payload.scope_id, + administrative=has_scope(principal, ADMIN_SCOPE), + ) + scope_id = validate_access_scope_reference( + get_registry(), + tenant_id=tenant_id or principal.tenant_id, + scope_type=scope_type, + scope_id=scope_id, + preserve_existing=( + existing.scope_id + if existing.scope_type == scope_type + else None + ), + ) + payload = payload.model_copy( + update={"scope_type": scope_type, "scope_id": scope_id} + ) pipeline = update_pipeline( session, tenant_id=principal.tenant_id, @@ -495,7 +562,7 @@ def api_update_pipeline( actor_id=_actor_id(principal), payload=payload, ) - except PermissionError as exc: + except (PermissionError, ValueError) as exc: raise _governance_http_error(exc) from exc except DataflowError as exc: raise _http_error(exc) from exc @@ -579,6 +646,13 @@ def api_derive_pipeline( payload = payload.model_copy( update={"scope_type": scope_type, "scope_id": scope_id} ) + scope_id = validate_access_scope_reference( + get_registry(), + tenant_id=tenant_id or principal.tenant_id, + scope_type=scope_type, + scope_id=scope_id, + ) + payload = payload.model_copy(update={"scope_id": scope_id}) pipeline = derive_pipeline( session, tenant_id=tenant_id or principal.tenant_id, diff --git a/src/govoplan_dataflow/backend/triggers.py b/src/govoplan_dataflow/backend/triggers.py index 111ebc7..bd90493 100644 --- a/src/govoplan_dataflow/backend/triggers.py +++ b/src/govoplan_dataflow/backend/triggers.py @@ -707,7 +707,15 @@ def _dispatch_delivery( registry: object | None, now: datetime, ) -> str: - trigger = session.get(DataflowTrigger, delivery.trigger_id) + # Serialize capacity checks per trigger. Delivery rows are independently + # leased with SKIP LOCKED, so without this lock two workers could each + # claim a different delivery and both observe the same available slot. + trigger = session.scalar( + select(DataflowTrigger) + .where(DataflowTrigger.id == delivery.trigger_id) + .with_for_update() + .execution_options(populate_existing=True) + ) pipeline = session.get(DataflowPipeline, delivery.pipeline_id) if trigger is None or pipeline is None: return _finish_delivery( diff --git a/tests/test_service.py b/tests/test_service.py index 5d2bbf6..bb77fc8 100644 --- a/tests/test_service.py +++ b/tests/test_service.py @@ -23,6 +23,7 @@ from govoplan_dataflow.backend.db.models import ( DataflowPipelineRevision, DataflowRun, ) +from govoplan_dataflow.backend.governance import normalize_definition_scope from govoplan_dataflow.backend.schemas import ( GraphEdge, GraphNode, @@ -140,6 +141,18 @@ class FakeRegistry: class DataflowServiceTests(unittest.TestCase): + def test_user_scope_normalizes_membership_to_account_id(self) -> None: + tenant_id, scope_type, scope_id = normalize_definition_scope( + principal(), + scope_type="user", + scope_id="membership-1", + administrative=False, + ) + + self.assertEqual("tenant-1", tenant_id) + self.assertEqual("user", scope_type) + self.assertEqual("account-1", scope_id) + def setUp(self) -> None: self.engine = create_engine("sqlite:///:memory:") Base.metadata.create_all( diff --git a/webui/package.json b/webui/package.json index 4aa8636..dbe509c 100644 --- a/webui/package.json +++ b/webui/package.json @@ -23,7 +23,7 @@ "lucide-react": "^1.23.0", "react": "^19.0.0", "react-dom": "^19.0.0", - "react-router-dom": "^7.1.1", + "react-router-dom": ">=7.18.2 <8", "typescript": "^5.7.2" }, "peerDependenciesMeta": { diff --git a/webui/src/api/dataflow.ts b/webui/src/api/dataflow.ts index b5d1ba7..6d3ae18 100644 --- a/webui/src/api/dataflow.ts +++ b/webui/src/api/dataflow.ts @@ -1,4 +1,9 @@ -import { apiFetch, type ApiSettings } from "@govoplan/core-webui"; +import { + apiFetch, + apiReferenceOptionProvider, + type ApiSettings, + type ReferenceOptionProvider +} from "@govoplan/core-webui"; import type { DefinitionGraphNodeType } from "@govoplan/core-webui/definition-graph"; export type PipelineStatus = "draft" | "active" | "archived"; @@ -376,6 +381,17 @@ export function deriveDataflowPipeline( ); } +export function dataflowScopeReferenceProvider( + settings: ApiSettings, + scopeType: "user" | "group" +): ReferenceOptionProvider { + return apiReferenceOptionProvider( + settings, + "/api/v1/dataflow/scope-targets", + { scope_type: scopeType } + ); +} + export async function listDataflowTriggers( settings: ApiSettings, pipelineId: string diff --git a/webui/src/features/dataflow/DataflowPage.tsx b/webui/src/features/dataflow/DataflowPage.tsx index 1287fe7..c9b28ab 100644 --- a/webui/src/features/dataflow/DataflowPage.tsx +++ b/webui/src/features/dataflow/DataflowPage.tsx @@ -31,6 +31,7 @@ import { FormField, IconButton, LoadingFrame, + ReferenceSelect, SegmentedControl, StatusBadge, ToggleSwitch, @@ -49,6 +50,7 @@ import { createDataflowSourceSnapshot, deleteDataflowPipeline, deleteDataflowTrigger, + dataflowScopeReferenceProvider, deriveDataflowPipeline, listDataflowNodeTypes, listDataflowPipelineRuns, @@ -939,6 +941,7 @@ export default function DataflowPage({ settings, auth }: { settings: ApiSettings /> ) => void; onClose: () => void; }) { const provenance = draft?.governance?.actions.edit?.source_path ?? []; + const referenceScopeType = draft?.scopeType === "group" ? "group" : "user"; + const scopeProvider = useMemo( + () => dataflowScopeReferenceProvider(settings, referenceScopeType), + [ + referenceScopeType, + settings.accessToken, + settings.apiBaseUrl, + settings.apiKey + ] + ); return ( User - - onChange({ scopeId: event.target.value })} - /> - + > + onChange({ scopeId: value })} + provider={scopeProvider} + aria-label={ + draft.scopeType === "user" + ? "Definition user" + : "Definition group" + } + placeholder={ + draft.scopeType === "user" + ? "Select a user" + : "Select a group" + } + disabled={!editable || Boolean(draft.id)} + required + /> + + ) : null} + dataflowScopeReferenceProvider( + settings, + scopeType === "group" ? "group" : "user" + ), + [ + scopeType, + settings.accessToken, + settings.apiBaseUrl, + settings.apiKey + ] + ); useEffect(() => { if (!open) return; @@ -1177,7 +1216,12 @@ function DerivePipelineDialog({ @@ -1192,7 +1236,10 @@ function DerivePipelineDialog({ {scopeType !== "tenant" ? ( - - + setScopeId(event.target.value)} + onChange={(value) => setScopeId(value)} + provider={scopeProvider} + aria-label={ + scopeType === "user" ? "Copy target user" : "Copy target group" + } + placeholder={ + scopeType === "user" ? "Select a user" : "Select a group" + } + disabled={busy} + required /> ) : null} diff --git a/webui/src/features/dataflow/DataflowPipelinesWidget.tsx b/webui/src/features/dataflow/DataflowPipelinesWidget.tsx new file mode 100644 index 0000000..ac6b67a --- /dev/null +++ b/webui/src/features/dataflow/DataflowPipelinesWidget.tsx @@ -0,0 +1,103 @@ +import { useCallback } from "react"; +import { Waypoints } from "lucide-react"; +import { Link } from "react-router-dom"; +import { + DashboardWidgetList, + DismissibleAlert, + LoadingFrame, + StatusBadge, + useDashboardWidgetData, + type ApiSettings, + type DashboardWidgetConfiguration +} from "@govoplan/core-webui"; +import { + listDataflowPipelines, + type Pipeline +} from "../../api/dataflow"; + +export default function DataflowPipelinesWidget({ + settings, + refreshKey, + configuration +}: { + settings: ApiSettings; + refreshKey: number; + configuration: DashboardWidgetConfiguration; +}) { + const maxItems = numberSetting(configuration.maxItems, 5, 1, 12); + const includeDrafts = configuration.includeDrafts !== false; + const load = useCallback(async () => { + const pipelines = await listDataflowPipelines(settings); + return pipelines + .filter((pipeline) => includeDrafts || pipeline.status !== "draft") + .sort(comparePipelines) + .slice(0, maxItems); + }, [includeDrafts, maxItems, settings]); + const { data: pipelines, loading, error } = useDashboardWidgetData( + load, + refreshKey + ); + + return ( + + {error && ( + + {error} + + )} + ({ + id: pipeline.id, + title: pipeline.name, + detail: pipelineDescription(pipeline), + meta: scopeLabel(pipeline), + leading: + ); +} + +function comparePipelines(left: Pipeline, right: Pipeline): number { + if (left.status !== right.status) { + if (left.status === "active") return -1; + if (right.status === "active") return 1; + } + return ( + new Date(right.updated_at).getTime() - new Date(left.updated_at).getTime() + ); +} + +function pipelineDescription(pipeline: Pipeline): string { + const nodeCount = pipeline.revision.graph.nodes.length; + const kind = + pipeline.governance.definition_kind === "template" ? "template" : "flow"; + return `${nodeCount} ${nodeCount === 1 ? "node" : "nodes"} ยท ${kind}`; +} + +function scopeLabel(pipeline: Pipeline): string { + const scope = pipeline.governance.scope_type; + return scope.charAt(0).toUpperCase() + scope.slice(1); +} + +function numberSetting( + value: unknown, + fallback: number, + minimum: number, + maximum: number +): number { + const numeric = typeof value === "number" ? value : Number(value); + return Number.isFinite(numeric) + ? Math.max(minimum, Math.min(maximum, Math.floor(numeric))) + : fallback; +} diff --git a/webui/src/module.ts b/webui/src/module.ts index ff30677..d4f3bf9 100644 --- a/webui/src/module.ts +++ b/webui/src/module.ts @@ -1,11 +1,59 @@ import { createElement, lazy } from "react"; -import type { PlatformWebModule } from "@govoplan/core-webui"; +import type { + DashboardWidgetsUiCapability, + PlatformWebModule +} from "@govoplan/core-webui"; +import DataflowPipelinesWidget from "./features/dataflow/DataflowPipelinesWidget"; import "@xyflow/react/dist/style.css"; import "./styles/dataflow.css"; const DataflowPage = lazy(() => import("./features/dataflow/DataflowPage")); const readScopes = ["dataflow:pipeline:read", "dataflow:pipeline:admin"]; +const dataflowDashboardWidgets: DashboardWidgetsUiCapability = { + widgets: [ + { + id: "dataflow.pipelines", + surfaceId: "dataflow.widget.pipelines", + title: "Dataflows", + description: "Active and recently edited dataflow definitions.", + moduleId: "dataflow", + category: "Automation", + order: 70, + defaultVisible: false, + defaultSize: "medium", + supportedSizes: ["medium", "wide"], + anyOf: readScopes, + refreshIntervalMs: 60_000, + defaultConfiguration: { + maxItems: 5, + includeDrafts: true + }, + configurationFields: [ + { + id: "maxItems", + label: "Maximum dataflows", + kind: "number", + min: 1, + max: 12, + step: 1, + required: true + }, + { + id: "includeDrafts", + label: "Include drafts", + kind: "boolean" + } + ], + render: ({ settings, refreshKey, configuration }) => + createElement(DataflowPipelinesWidget, { + settings, + refreshKey, + configuration + }) + } + ] +}; export const dataflowModule: PlatformWebModule = { id: "dataflow", @@ -22,6 +70,15 @@ export const dataflowModule: PlatformWebModule = { "risk_compliance", "workflow" ], + viewSurfaces: [ + { + id: "dataflow.widget.pipelines", + moduleId: "dataflow", + kind: "section", + label: "Dataflows widget", + order: 70 + } + ], navItems: [ { to: "/dataflow", label: "Dataflow", iconName: "waypoints", anyOf: readScopes, order: 72 } ], @@ -32,7 +89,10 @@ export const dataflowModule: PlatformWebModule = { order: 72, render: ({ settings, auth }) => createElement(DataflowPage, { settings, auth }) } - ] + ], + uiCapabilities: { + "dashboard.widgets": dataflowDashboardWidgets + } }; export default dataflowModule;