From cdc8aa4bf7bcdaba9253edcf38e84b69c4621385 Mon Sep 17 00:00:00 2001 From: mo Date: Sat, 27 Jun 2026 21:10:16 +0000 Subject: [PATCH] feat(pii): masking for Cassandra & Neo4j (graph nodes + live property read) + MongoDB payload/free-text --- api/dataflow.py | 16 ++++++++++++---- api/pii_catalog.py | 30 ++++++++++++++++++++++++++++-- 2 files changed, 40 insertions(+), 6 deletions(-) diff --git a/api/dataflow.py b/api/dataflow.py index 15778aa..e29a19f 100644 --- a/api/dataflow.py +++ b/api/dataflow.py @@ -34,9 +34,11 @@ NODES: list[dict[str, Any]] = [ {"id": "generator", "label": "Data Generator", "sub": "Airflow DAGs", "kind": "generator", "x": 9, "y": 34}, {"id": "hdfs", "label": "Hadoop HDFS", "sub": "historical_sales", "kind": "hadoop", "x": 9, "y": 80}, # col 1 — source databases - {"id": "postgres", "label": "PostgreSQL", "sub": "sales_orders", "kind": "source", "x": 28, "y": 16}, - {"id": "mysql", "label": "MySQL", "sub": "employee_events", "kind": "source", "x": 28, "y": 38}, - {"id": "mongodb", "label": "MongoDB", "sub": "events", "kind": "source", "x": 28, "y": 60}, + {"id": "postgres", "label": "PostgreSQL", "sub": "sales_orders", "kind": "source", "x": 28, "y": 12}, + {"id": "mysql", "label": "MySQL", "sub": "employee_events", "kind": "source", "x": 28, "y": 28}, + {"id": "mongodb", "label": "MongoDB", "sub": "events", "kind": "source", "x": 28, "y": 44}, + {"id": "cassandra", "label": "Cassandra", "sub": "device_metrics", "kind": "source", "x": 28, "y": 60}, + {"id": "neo4j", "label": "Neo4j", "sub": "Product · Supplier graph", "kind": "source", "x": 28, "y": 76}, # col 2 — change data capture {"id": "kafka", "label": "Kafka · Debezium", "sub": "CDC topics", "kind": "stream", "x": 47, "y": 24}, {"id": "spark", "label": "Apache Spark", "sub": "Streaming · batch", "kind": "compute", "x": 62, "y": 38, @@ -57,9 +59,13 @@ EDGES: list[dict[str, Any]] = [ {"from": "generator", "to": "postgres", "kind": "generate", "movement_id": "gen_postgres"}, {"from": "generator", "to": "mysql", "kind": "generate", "movement_id": "gen_mysql"}, {"from": "generator", "to": "mongodb", "kind": "generate", "movement_id": "gen_mongodb"}, + {"from": "generator", "to": "cassandra", "kind": "generate"}, + {"from": "generator", "to": "neo4j", "kind": "generate"}, {"from": "postgres", "to": "kafka", "kind": "cdc"}, {"from": "mysql", "to": "kafka", "kind": "cdc"}, {"from": "mongodb", "to": "kafka", "kind": "cdc"}, + {"from": "cassandra", "to": "kafka", "kind": "cdc"}, + {"from": "neo4j", "to": "kafka", "kind": "cdc"}, {"from": "kafka", "to": "spark", "kind": "stream"}, {"from": "spark", "to": "iceberg_curated", "kind": "movement", "movement_id": "spark_to_curated"}, {"from": "spark", "to": "s3_cdc", "kind": "movement", "movement_id": "spark_to_s3"}, @@ -77,6 +83,8 @@ EDGES: list[dict[str, Any]] = [ {"from": "postgres", "to": "openmetadata", "kind": "catalog"}, {"from": "mysql", "to": "openmetadata", "kind": "catalog"}, {"from": "mongodb", "to": "openmetadata", "kind": "catalog"}, + {"from": "cassandra", "to": "openmetadata", "kind": "catalog"}, + {"from": "neo4j", "to": "openmetadata", "kind": "catalog"}, {"from": "trino", "to": "openmetadata", "kind": "catalog"}, ] @@ -165,7 +173,7 @@ async def _build() -> dict[str, Any]: used = spark.get("cores_used") or 0 metric = f"{spark.get('alive_workers', 0)} workers · {used}/{cores} cores · {apps} apps" node["level"] = "ok" if spark.get("ui_ok") and (spark.get("status") or "").upper() == "ALIVE" else "warn" - elif n["id"] in ("postgres", "mysql", "mongodb"): + elif n["id"] in ("postgres", "mysql", "mongodb", "cassandra", "neo4j"): metric = f"{cdc.get('by_source', {}).get(n['id'], 0)} CDC/15m" elif n["id"] == "iceberg_hadoop": c = _iceberg_hadoop_count() diff --git a/api/pii_catalog.py b/api/pii_catalog.py index 4fc5e22..d5a149e 100644 --- a/api/pii_catalog.py +++ b/api/pii_catalog.py @@ -62,6 +62,14 @@ DATASETS = [ "table": "iceberg.hadoop.historical_sales_hdfs", "table_name": "historical_sales_hdfs", "catalog": "iceberg", "schema": "hadoop", "om_fqn": "atc_trino.iceberg.hadoop.historical_sales_hdfs"}, + {"key": "cassandra", "node_id": "cassandra", "label": "Cassandra device_metrics", + "table": "cassandra_telemetry.telemetry.device_metrics", "table_name": "device_metrics", + "catalog": "cassandra_telemetry", "schema": "telemetry", + "om_fqn": "atc_trino.cassandra_telemetry.telemetry.device_metrics", + "native": {"engine": "cassandra"}}, + {"key": "neo4j", "node_id": "neo4j", "label": "Neo4j Product/Supplier graph", + "table": "neo4j_graph", "table_name": "graph", "catalog": "neo4j", + "native": {"engine": "neo4j"}}, ] KEY_BY_NODE = {ds["node_id"]: ds["key"] for ds in DATASETS} @@ -112,7 +120,8 @@ PII_RULES: list[tuple[str, str]] = [ (r"address|street|city|zip|postal|postcode", "ADDRESS"), (r"dob|birth|date_of_birth", "DOB"), (r"ip_addr|ip_address|_ip$|customer_ip|client_ip", "IP"), - (r"customer_id|client_id|user_id|account_id|member_id|device_id|subscriber_id|employee_id|guest_id|person_id", "IDENTIFIER"), + (r"customer_id|client_id|user_id|account_id|member_id|device_id|subscriber_id|employee_id|guest_id|person_id|supplier_id", "IDENTIFIER"), + (r"^payload$|payload|^notes$|^note$|free_text|freeform|raw_json|description", "FREEFORM"), ] _cache: dict[str, Any] = {"ts": 0.0, "data": None} @@ -182,13 +191,30 @@ def _trino_columns(catalog: str, schema: str | None, table_name: str) -> list[st return [] +def _neo4j_property_keys() -> list[str]: + """Property keys across the Neo4j graph, used as the 'columns' for PII tagging.""" + try: + from neo4j import GraphDatabase + uri = os.getenv("NEO4J_URI", f"bolt://{os.getenv('DB_HOST', '10.0.21.51')}:7687") + drv = GraphDatabase.driver(uri, auth=(os.getenv("NEO4J_USER", "neo4j"), os.getenv("NEO4J_PASSWORD", "testpwd"))) + with drv.session() as s: + keys = [r["propertyKey"] for r in s.run("CALL db.propertyKeys() YIELD propertyKey RETURN propertyKey")] + drv.close() + return keys + except Exception: + return [] + + def _build() -> dict[str, Any]: datasets_out = [] total_pii = 0 total_masked = 0 om_used = False for ds in DATASETS: - cols = _trino_columns(ds["catalog"], ds.get("schema"), ds["table_name"]) + if (ds.get("native") or {}).get("engine") == "neo4j": + cols = _neo4j_property_keys() + else: + cols = _trino_columns(ds["catalog"], ds.get("schema"), ds["table_name"]) om_tags = _om_column_tags(ds["om_fqn"]) if ds.get("om_fqn") else {} if om_tags: om_used = True