From 3ea5ee4fa6ea96ddc0bb9471fec47d9a8125f65d Mon Sep 17 00:00:00 2001 From: ntkathole Date: Mon, 14 Sep 2026 09:53:31 +0530 Subject: [PATCH] feat: OpenLineage consumer improvements - dataset versioning, column lineage, ownership Signed-off-by: ntkathole --- sdk/python/feast/openlineage/consumer.py | 54 ++ sdk/python/feast/openlineage/emitter.py | 139 ++++- sdk/python/feast/openlineage/models.py | 121 +++++ sdk/python/feast/openlineage/processor.py | 132 ++++- sdk/python/feast/openlineage/store.py | 481 ++++++++++++++++++ .../tests/unit/openlineage/test_processor.py | 385 ++++++++++++++ .../tests/unit/openlineage/test_store.py | 291 +++++++++++ ui/src/components/OpenLineageGraph.tsx | 15 + ui/src/queries/useLoadOpenLineageGraph.ts | 4 + 9 files changed, 1615 insertions(+), 7 deletions(-) diff --git a/sdk/python/feast/openlineage/consumer.py b/sdk/python/feast/openlineage/consumer.py index 886fc680baa..b8489ef30db 100644 --- a/sdk/python/feast/openlineage/consumer.py +++ b/sdk/python/feast/openlineage/consumer.py @@ -337,6 +337,10 @@ def get_full_lineage_graph( "schema": schema, "description": ds.get("description"), "source_type": ds.get("source_type"), + "owner": ds.get("owner_name"), + "owner_type": ds.get("owner_type"), + "lifecycle_state": ds.get("lifecycle_state"), + "current_version": ds.get("current_version"), "facets": facets, } ) @@ -377,6 +381,56 @@ def get_full_lineage_graph( "total_nodes": total_nodes, } + # ── Dataset versioning endpoints ── + + @router.get("/lineage/openlineage/datasets/{namespace}/{name}/versions") + def list_dataset_versions( + namespace: str, + name: str, + limit: int = Query(50, ge=1, le=500), + offset: int = Query(0, ge=0), + ): + """List version history for a dataset.""" + ns_filter = _get_namespace_filter(get_allowed_namespaces) + if ns_filter is not None and namespace not in ns_filter: + return {"versions": [], "total": 0} + versions = store.get_dataset_versions(namespace, name, limit, offset) + for v in versions: + v["schema"] = _safe_parse_json(v.pop("schema_json", None)) + v["facets"] = _safe_parse_json(v.pop("facets_json", None)) + return {"versions": versions, "total": len(versions)} + + # ── Column-level lineage endpoint ── + + @router.get("/lineage/openlineage/datasets/{namespace}/{name}/columns") + def get_column_lineage( + namespace: str, + name: str, + direction: str = Query("both", pattern="^(both|upstream|downstream)$"), + ): + """Get column-level lineage for a dataset. + + Shows which input fields map to which output fields, including + transformation types when available. + """ + ns_filter = _get_namespace_filter(get_allowed_namespaces) + if ns_filter is not None and namespace not in ns_filter: + return {"column_lineage": []} + cl = store.get_column_lineage(namespace, name, direction) + return {"column_lineage": cl} + + # ── Execution tree endpoint ── + + @router.get("/lineage/openlineage/runs/{run_id}/tree") + def get_run_tree(run_id: str): + """Get the full execution tree (all runs sharing the same root). + + Each run includes parent_run_id and root_run_id for hierarchy + reconstruction on the client side. + """ + tree = store.get_run_tree(run_id) + return {"runs": tree, "total": len(tree)} + # ── Run history endpoints ── @router.get("/lineage/openlineage/runs") diff --git a/sdk/python/feast/openlineage/emitter.py b/sdk/python/feast/openlineage/emitter.py index a93eca7f23d..f1fdcb911be 100644 --- a/sdk/python/feast/openlineage/emitter.py +++ b/sdk/python/feast/openlineage/emitter.py @@ -567,10 +567,13 @@ def emit_saved_dataset_lineage( """ Emit lineage for a saved dataset definition. - Creates a definition job with inputs derived from: - - FeatureService (when feature_service_name is set) - - FeatureViews (extracted from feature refs in the format "view:feat") - - DataSources (matched by comparing storage location against registered sources) + Emits two events: + 1. A **RunEvent** that captures the definition-time topology: + FeatureService/FeatureView → SavedDataset job. + 2. A **DatasetEvent** with SymlinksDatasetFacet that links the + logical SavedDataset name to its physical storage URI, plus + OwnershipDatasetFacet and LifecycleStateChangeDatasetFacet + when the information is available. Args: saved_dataset: The SavedDataset object @@ -675,7 +678,7 @@ def emit_saved_dataset_lineage( **self._job_kind_facets(FeastJobKind.DEFINITION, project), } - return self._client.emit_run_event( + run_result = self._client.emit_run_event( job_name=f"saved_dataset_{saved_dataset.name}", run_id=str(uuid.uuid4()), event_type=RunState.COMPLETE, @@ -684,12 +687,138 @@ def emit_saved_dataset_lineage( job_facets=job_facets, namespace=namespace, ) + + # Emit a DatasetEvent with symlink, ownership, and lifecycle facets + self._emit_saved_dataset_metadata( + saved_dataset, project, namespace, registered_data_sources + ) + + return run_result except Exception as e: logger.error( f"Error emitting saved dataset lineage for {saved_dataset.name}: {e}" ) return False + def _emit_saved_dataset_metadata( + self, + saved_dataset: Any, + project: str, + namespace: str, + registered_data_sources: Optional[List[Any]] = None, + ): + """Emit a DatasetEvent for a SavedDataset with OL-standard facets. + + Includes: + - SymlinksDatasetFacet linking logical name to physical storage URI + - OwnershipDatasetFacet from tags (``owner`` key) or the object owner + - LifecycleStateChangeDatasetFacet (CREATE) + - SchemaDatasetFacet from features + """ + try: + from openlineage.client.facet_v2 import ( + lifecycle_state_change_dataset, + ownership_dataset, + schema_dataset, + symlinks_dataset, + ) + except ImportError: + logger.debug( + "openlineage facet_v2 modules not available; " + "skipping DatasetEvent for SavedDataset" + ) + return + + ds_facets: Dict[str, Any] = {} + + # SymlinksDatasetFacet: link logical name → physical storage URI + storage_uri = self._get_saved_dataset_storage_uri(saved_dataset) + if storage_uri: + ds_facets["symlinks"] = symlinks_dataset.SymlinksDatasetFacet( + identifiers=[ + symlinks_dataset.Identifier( + namespace=namespace, + name=storage_uri, + type="TABLE", + ) + ] + ) + + # OwnershipDatasetFacet + owner = None + if hasattr(saved_dataset, "tags") and saved_dataset.tags: + owner = saved_dataset.tags.get("owner") + if not owner and hasattr(saved_dataset, "owner") and saved_dataset.owner: + owner = saved_dataset.owner + if owner: + ds_facets["ownership"] = ownership_dataset.OwnershipDatasetFacet( + owners=[ + ownership_dataset.Owner( + name=owner, + type="FEAST_TAG", + ) + ] + ) + + # LifecycleStateChangeDatasetFacet + ds_facets["lifecycleStateChange"] = ( + lifecycle_state_change_dataset.LifecycleStateChangeDatasetFacet( + lifecycleStateChange=lifecycle_state_change_dataset.LifecycleStateChange.CREATE, + ) + ) + + # SchemaDatasetFacet + if saved_dataset.features: + ds_facets["schema"] = schema_dataset.SchemaDatasetFacet( + fields=[ + schema_dataset.SchemaDatasetFacetFields(name=f, type="UNKNOWN") + for f in saved_dataset.features + ] + ) + + try: + self._client.emit_dataset_event( + dataset_name=saved_dataset.name, + namespace=namespace, + facets=ds_facets, + ) + except Exception as e: + logger.warning( + f"Failed to emit DatasetEvent for SavedDataset " + f"{saved_dataset.name}: {e}" + ) + + @staticmethod + def _get_saved_dataset_storage_uri(saved_dataset: Any) -> Optional[str]: + """Extract a physical storage URI from a SavedDataset's storage.""" + if not hasattr(saved_dataset, "storage") or not saved_dataset.storage: + return None + try: + storage_proto = saved_dataset.storage.to_proto() + except Exception: + return None + + if hasattr(storage_proto, "file_url") and storage_proto.file_url: + return storage_proto.file_url + if hasattr(storage_proto, "bigquery_path") and storage_proto.bigquery_path: + return storage_proto.bigquery_path + if hasattr(storage_proto, "redshift_path") and storage_proto.redshift_path: + return storage_proto.redshift_path + if hasattr(storage_proto, "s3_path") and storage_proto.s3_path: + return storage_proto.s3_path + if hasattr(storage_proto, "trino_path") and storage_proto.trino_path: + return storage_proto.trino_path + + for attr_name in dir(storage_proto): + if attr_name.startswith("_"): + continue + if "path" in attr_name or "url" in attr_name or "uri" in attr_name: + val = getattr(storage_proto, attr_name, None) + if val and isinstance(val, str): + return val + + return None + def emit_materialize_start( self, feature_views: List["FeatureView"], diff --git a/sdk/python/feast/openlineage/models.py b/sdk/python/feast/openlineage/models.py index ec79cf3620d..195fe3177ee 100644 --- a/sdk/python/feast/openlineage/models.py +++ b/sdk/python/feast/openlineage/models.py @@ -17,12 +17,25 @@ Uses SQLAlchemy Core (Table + Column) pattern consistent with Feast's existing SQL registry in feast/infra/registry/sql.py. + +Tables +------ +Core: + openlineage_events, openlineage_jobs, openlineage_datasets, + openlineage_runs, openlineage_run_io, openlineage_lineage_edges, + openlineage_dataset_symlinks + +Extended (OL spec facets): + openlineage_dataset_versions – per-run snapshots of dataset state + openlineage_column_lineage – column-level field transformations + openlineage_dataset_ownership – extracted OwnershipDatasetFacet rows """ from sqlalchemy import ( BigInteger, Column, Index, + Integer, MetaData, String, Text, @@ -71,6 +84,11 @@ Column("feast_object_type", String(100), nullable=True), Column("feast_object_name", String(255), nullable=True), Column("feast_project", String(255), nullable=True), + # Extended: ownership, lifecycle, version tracking + Column("owner_name", String(512), nullable=True), + Column("owner_type", String(100), nullable=True), + Column("lifecycle_state", String(50), nullable=True), + Column("current_version", Integer, nullable=True), Column("updated_at", BigInteger, nullable=False), ) @@ -84,6 +102,9 @@ Column("started_at", BigInteger, nullable=True), Column("ended_at", BigInteger, nullable=True), Column("facets_json", Text, nullable=True), + # Extended: execution hierarchy from ParentRunFacet + Column("parent_run_id", String(255), nullable=True), + Column("root_run_id", String(255), nullable=True), Column("updated_at", BigInteger, nullable=False), ) @@ -145,6 +166,66 @@ ) +openlineage_dataset_versions = ( + "openlineage_dataset_versions", + ol_metadata, + Column("version_id", Integer, primary_key=True, autoincrement=True), + Column("dataset_namespace", String(512), nullable=False), + Column("dataset_name", String(512), nullable=False), + Column("version", Integer, nullable=False), + Column("created_by_run_id", String(255), nullable=True), + Column("schema_json", Text, nullable=True), + Column("facets_json", Text, nullable=True), + Column("created_at", BigInteger, nullable=False), + UniqueConstraint( + "dataset_namespace", + "dataset_name", + "version", + name="uq_dataset_version", + ), +) + +openlineage_column_lineage = ( + "openlineage_column_lineage", + ol_metadata, + Column("dataset_namespace", String(512), nullable=False), + Column("dataset_name", String(512), nullable=False), + Column("output_field", String(512), nullable=False), + Column("input_namespace", String(512), nullable=False), + Column("input_name", String(512), nullable=False), + Column("input_field", String(512), nullable=False), + Column("transformation_type", String(100), nullable=True), + Column("transformation_description", Text, nullable=True), + Column("run_id", String(255), nullable=True), + Column("updated_at", BigInteger, nullable=False), + UniqueConstraint( + "dataset_namespace", + "dataset_name", + "output_field", + "input_namespace", + "input_name", + "input_field", + name="uq_column_lineage", + ), +) + +openlineage_dataset_ownership = ( + "openlineage_dataset_ownership", + ol_metadata, + Column("dataset_namespace", String(512), nullable=False), + Column("dataset_name", String(512), nullable=False), + Column("owner_name", String(512), nullable=False), + Column("owner_type", String(100), nullable=True), + Column("updated_at", BigInteger, nullable=False), + UniqueConstraint( + "dataset_namespace", + "dataset_name", + "owner_name", + name="uq_dataset_owner", + ), +) + + def _build_tables(): """Build Table objects from the tuple definitions above.""" from sqlalchemy import Table @@ -158,6 +239,9 @@ def _build_tables(): ("run_io", openlineage_run_io), ("lineage_edges", openlineage_lineage_edges), ("dataset_symlinks", openlineage_dataset_symlinks), + ("dataset_versions", openlineage_dataset_versions), + ("column_lineage", openlineage_column_lineage), + ("dataset_ownership", openlineage_dataset_ownership), ]: tbl_name = tbl_def[0] meta = tbl_def[1] @@ -216,3 +300,40 @@ def _build_tables(): OL_TABLES["dataset_symlinks"].c.linked_namespace, OL_TABLES["dataset_symlinks"].c.linked_name, ) + +idx_dataset_versions_ds = Index( + "idx_ol_dv_dataset", + OL_TABLES["dataset_versions"].c.dataset_namespace, + OL_TABLES["dataset_versions"].c.dataset_name, +) +idx_dataset_versions_run = Index( + "idx_ol_dv_run", + OL_TABLES["dataset_versions"].c.created_by_run_id, +) +idx_column_lineage_ds = Index( + "idx_ol_cl_dataset", + OL_TABLES["column_lineage"].c.dataset_namespace, + OL_TABLES["column_lineage"].c.dataset_name, +) +idx_column_lineage_input = Index( + "idx_ol_cl_input", + OL_TABLES["column_lineage"].c.input_namespace, + OL_TABLES["column_lineage"].c.input_name, +) +idx_dataset_ownership_ds = Index( + "idx_ol_do_dataset", + OL_TABLES["dataset_ownership"].c.dataset_namespace, + OL_TABLES["dataset_ownership"].c.dataset_name, +) +idx_runs_parent = Index( + "idx_ol_runs_parent", + OL_TABLES["runs"].c.parent_run_id, +) +idx_runs_root = Index( + "idx_ol_runs_root", + OL_TABLES["runs"].c.root_run_id, +) +idx_datasets_lifecycle = Index( + "idx_ol_datasets_lifecycle", + OL_TABLES["datasets"].c.lifecycle_state, +) diff --git a/sdk/python/feast/openlineage/processor.py b/sdk/python/feast/openlineage/processor.py index de1c4843969..3e5f4ea185e 100644 --- a/sdk/python/feast/openlineage/processor.py +++ b/sdk/python/feast/openlineage/processor.py @@ -19,10 +19,11 @@ extracts metadata, and builds the lineage graph in the store. """ +import json import logging import re import uuid -from typing import Any, Dict, List, Optional +from typing import Any, Dict, List, Optional, Tuple from urllib.parse import urlparse, urlunparse from feast.openlineage.store import OpenLineageStore @@ -158,8 +159,18 @@ def _process_run_event(self, event_id: str, event: Dict[str, Any]): self._store.upsert_job(job_namespace, job_name, job, producer=producer) + # Extract execution hierarchy from ParentRunFacet run_facets = run.get("facets", {}) - self._store.upsert_run(run_id, job_namespace, job_name, event_type, run_facets) + parent_run_id, root_run_id = self._extract_parent_hierarchy(run_facets, run_id) + self._store.upsert_run( + run_id, + job_namespace, + job_name, + event_type, + run_facets, + parent_run_id=parent_run_id, + root_run_id=root_run_id, + ) # Parent/child job link (e.g. Feast materialize → SparkApplication) parent = run_facets.get("parent") or run_facets.get("parentRun") @@ -194,6 +205,7 @@ def _process_run_event(self, event_id: str, event: Dict[str, Any]): ) self._store.store_run_io(run_id, ds_namespace, ds_name, "INPUT", ds_facets) self._process_dataset_symlinks(ds_namespace, ds_name, ds_facets) + self._process_column_lineage(ds_namespace, ds_name, ds_facets, run_id) self._store.upsert_lineage_edge( source_type="dataset", @@ -218,6 +230,7 @@ def _process_run_event(self, event_id: str, event: Dict[str, Any]): ) self._store.store_run_io(run_id, ds_namespace, ds_name, "OUTPUT", ds_facets) self._process_dataset_symlinks(ds_namespace, ds_name, ds_facets) + self._process_column_lineage(ds_namespace, ds_name, ds_facets, run_id) self._store.upsert_lineage_edge( source_type="job", @@ -229,6 +242,23 @@ def _process_run_event(self, event_id: str, event: Dict[str, Any]): edge_type="output", ) + # Create dataset version snapshot on COMPLETE events + if event_type == "COMPLETE" and run_id: + schema = ds_facets.get("schema") + try: + self._store.create_dataset_version( + namespace=ds_namespace, + name=ds_name, + run_id=run_id, + schema_json=json.dumps(schema) if schema else None, + facets_json=json.dumps(ds_facets) if ds_facets else None, + ) + except Exception as e: + logger.warning( + f"Failed to create dataset version for " + f"{ds_namespace}/{ds_name}: {e}" + ) + self._build_dataset_to_dataset_edges(inputs, outputs, job_namespace) def _process_dataset_event(self, event_id: str, event: Dict[str, Any]): @@ -436,6 +466,104 @@ def _process_dataset_symlinks( edge_type="symlink", ) + def _extract_parent_hierarchy( + self, + run_facets: Dict[str, Any], + run_id: str, + ) -> Tuple[Optional[str], Optional[str]]: + """Extract parent_run_id and root_run_id from ParentRunFacet. + + The OL spec defines ``parent`` (or ``parentRun``) as:: + + { "run": {"runId": "..."}, "job": {...}, + "_producer": "...", "_schemaURL": "..." } + + Returns (parent_run_id, root_run_id). When a ``root`` key is + present inside the parent facet, its ``runId`` is used as the + root; otherwise the parent is also the root. + """ + parent = run_facets.get("parent") or run_facets.get("parentRun") + if not isinstance(parent, dict): + return None, None + + p_run = parent.get("run") or {} + parent_run_id = p_run.get("runId") + + root = parent.get("root") or {} + root_run = root.get("run") or {} if isinstance(root, dict) else {} + root_run_id = root_run.get("runId") or parent_run_id + + return parent_run_id, root_run_id + + def _process_column_lineage( + self, + ds_namespace: str, + ds_name: str, + ds_facets: Dict[str, Any], + run_id: Optional[str] = None, + ): + """Extract and store ColumnLineageDatasetFacet entries. + + OL spec ``columnLineage``:: + + { "fields": { + "output_col": { + "inputFields": [ + { "namespace": "...", "name": "...", "field": "...", + "transformations": [{"type": "...", "description": "..."}] } + ] + } + }} + """ + cl = ds_facets.get("columnLineage", {}) + if not isinstance(cl, dict): + return + fields = cl.get("fields", {}) + if not isinstance(fields, dict): + return + + for output_field, field_info in fields.items(): + if not isinstance(field_info, dict): + continue + input_fields = field_info.get("inputFields", []) + if not isinstance(input_fields, list): + continue + for inp in input_fields: + if not isinstance(inp, dict): + continue + inp_ns = inp.get("namespace", ds_namespace) + inp_name = inp.get("name", "") + inp_field = inp.get("field", "") + if not inp_name or not inp_field: + continue + + xform_type = None + xform_desc = None + transformations = inp.get("transformations", []) + if isinstance(transformations, list) and transformations: + first_xform = transformations[0] + if isinstance(first_xform, dict): + xform_type = first_xform.get("type") + xform_desc = first_xform.get("description") + + try: + self._store.upsert_column_lineage( + dataset_namespace=ds_namespace, + dataset_name=ds_name, + output_field=output_field, + input_namespace=inp_ns, + input_name=inp_name, + input_field=inp_field, + transformation_type=xform_type, + transformation_description=xform_desc, + run_id=run_id, + ) + except Exception as e: + logger.warning( + f"Failed to store column lineage " + f"{ds_namespace}/{ds_name}.{output_field}: {e}" + ) + def _resolve_feast_mapping( self, namespace: str, diff --git a/sdk/python/feast/openlineage/store.py b/sdk/python/feast/openlineage/store.py index bf9ed43d78c..d6a02e4d8ea 100644 --- a/sdk/python/feast/openlineage/store.py +++ b/sdk/python/feast/openlineage/store.py @@ -177,6 +177,22 @@ def upsert_dataset( if "dataSource" in facets: source_type = facets["dataSource"].get("name") + # Extract ownership facet (first owner as primary) + owner_name = None + owner_type = None + ownership = facets.get("ownership", {}) + if isinstance(ownership, dict): + owners = ownership.get("owners", []) + if owners and isinstance(owners, list) and len(owners) > 0: + owner_name = owners[0].get("name") + owner_type = owners[0].get("type") + + # Extract lifecycle state + lifecycle_state = None + lsc = facets.get("lifecycleStateChange", {}) + if isinstance(lsc, dict) and lsc.get("lifecycleStateChange"): + lifecycle_state = lsc["lifecycleStateChange"] + feast_obj_type = feast_mapping.get("type") if feast_mapping else None feast_obj_name = feast_mapping.get("name") if feast_mapping else None feast_project = feast_mapping.get("project") if feast_mapping else None @@ -212,6 +228,12 @@ def upsert_dataset( values["description"] = description values["schema_json"] = schema_json + if owner_name: + values["owner_name"] = owner_name + values["owner_type"] = owner_type + if lifecycle_state: + values["lifecycle_state"] = lifecycle_state + if feast_obj_type and feast_obj_type != "unknown": values["feast_object_type"] = feast_obj_type if feast_obj_name: @@ -238,6 +260,12 @@ def upsert_dataset( values["dataset_name"] = name conn.execute(tbl.insert().values(**values)) + # Persist all owners from OwnershipDatasetFacet + if isinstance(ownership, dict): + owners = ownership.get("owners", []) + if owners: + self._upsert_dataset_owners(namespace, name, owners) + def upsert_run( self, run_id: str, @@ -245,6 +273,8 @@ def upsert_run( job_name: str, state: str, facets: Optional[Dict] = None, + parent_run_id: Optional[str] = None, + root_run_id: Optional[str] = None, ): now = int(time.time() * 1000) tbl = OL_TABLES["runs"] @@ -258,6 +288,10 @@ def upsert_run( update_vals["ended_at"] = now if facets: update_vals["facets_json"] = json.dumps(facets) + if parent_run_id: + update_vals["parent_run_id"] = parent_run_id + if root_run_id: + update_vals["root_run_id"] = root_run_id conn.execute( tbl.update().where(tbl.c.run_id == run_id).values(**update_vals) ) @@ -273,6 +307,8 @@ def upsert_run( if state in ("COMPLETE", "FAIL", "ABORT") else None, facets_json=json.dumps(facets) if facets else None, + parent_run_id=parent_run_id, + root_run_id=root_run_id, updated_at=now, ) ) @@ -717,6 +753,395 @@ def get_all_symlinks(self) -> List[Dict[str, Any]]: rows = conn.execute(select(tbl)).fetchall() return [dict(r._mapping) for r in rows] + # ── Dataset versioning ── + + def create_dataset_version( + self, + namespace: str, + name: str, + run_id: Optional[str] = None, + schema_json: Optional[str] = None, + facets_json: Optional[str] = None, + ) -> int: + """Create a new version snapshot for a dataset. + + Returns the new version number. + """ + now = int(time.time() * 1000) + tbl = OL_TABLES["dataset_versions"] + tbl_ds = OL_TABLES["datasets"] + + with self._engine.begin() as conn: + max_ver = conn.execute( + select(func.max(tbl.c.version)).where( + tbl.c.dataset_namespace == namespace, + tbl.c.dataset_name == name, + ) + ).scalar() + new_version = (max_ver or 0) + 1 + + conn.execute( + tbl.insert().values( + dataset_namespace=namespace, + dataset_name=name, + version=new_version, + created_by_run_id=run_id, + schema_json=schema_json, + facets_json=facets_json, + created_at=now, + ) + ) + + conn.execute( + tbl_ds.update() + .where( + tbl_ds.c.dataset_namespace == namespace, + tbl_ds.c.dataset_name == name, + ) + .values(current_version=new_version, updated_at=now) + ) + + return new_version + + def get_dataset_versions( + self, + namespace: str, + name: str, + limit: int = 50, + offset: int = 0, + ) -> List[Dict[str, Any]]: + """Get version history for a dataset.""" + tbl = OL_TABLES["dataset_versions"] + query = ( + select(tbl) + .where( + tbl.c.dataset_namespace == namespace, + tbl.c.dataset_name == name, + ) + .order_by(tbl.c.version.desc()) + .limit(limit) + .offset(offset) + ) + with self._engine.connect() as conn: + rows = conn.execute(query).fetchall() + return [dict(r._mapping) for r in rows] + + def get_dataset_version( + self, + namespace: str, + name: str, + version: int, + ) -> Optional[Dict[str, Any]]: + """Get a specific dataset version.""" + tbl = OL_TABLES["dataset_versions"] + with self._engine.connect() as conn: + row = conn.execute( + select(tbl).where( + tbl.c.dataset_namespace == namespace, + tbl.c.dataset_name == name, + tbl.c.version == version, + ) + ).first() + return dict(row._mapping) if row else None + + # ── Column-level lineage ── + + def upsert_column_lineage( + self, + dataset_namespace: str, + dataset_name: str, + output_field: str, + input_namespace: str, + input_name: str, + input_field: str, + transformation_type: Optional[str] = None, + transformation_description: Optional[str] = None, + run_id: Optional[str] = None, + ): + """Store a column-level lineage mapping.""" + now = int(time.time() * 1000) + tbl = OL_TABLES["column_lineage"] + with self._engine.begin() as conn: + existing = conn.execute( + select(tbl).where( + tbl.c.dataset_namespace == dataset_namespace, + tbl.c.dataset_name == dataset_name, + tbl.c.output_field == output_field, + tbl.c.input_namespace == input_namespace, + tbl.c.input_name == input_name, + tbl.c.input_field == input_field, + ) + ).first() + + if existing: + conn.execute( + tbl.update() + .where( + tbl.c.dataset_namespace == dataset_namespace, + tbl.c.dataset_name == dataset_name, + tbl.c.output_field == output_field, + tbl.c.input_namespace == input_namespace, + tbl.c.input_name == input_name, + tbl.c.input_field == input_field, + ) + .values( + transformation_type=transformation_type, + transformation_description=transformation_description, + run_id=run_id, + updated_at=now, + ) + ) + else: + conn.execute( + tbl.insert().values( + dataset_namespace=dataset_namespace, + dataset_name=dataset_name, + output_field=output_field, + input_namespace=input_namespace, + input_name=input_name, + input_field=input_field, + transformation_type=transformation_type, + transformation_description=transformation_description, + run_id=run_id, + updated_at=now, + ) + ) + + def get_column_lineage( + self, + namespace: str, + name: str, + direction: str = "both", + ) -> List[Dict[str, Any]]: + """Get column-level lineage for a dataset. + + direction: "upstream" (inputs to this dataset), "downstream" + (where this dataset's fields go), or "both". + """ + tbl = OL_TABLES["column_lineage"] + results = [] + + with self._engine.connect() as conn: + if direction in ("both", "upstream"): + rows = conn.execute( + select(tbl).where( + tbl.c.dataset_namespace == namespace, + tbl.c.dataset_name == name, + ) + ).fetchall() + for r in rows: + entry = dict(r._mapping) + entry["direction"] = "upstream" + results.append(entry) + + if direction in ("both", "downstream"): + rows = conn.execute( + select(tbl).where( + tbl.c.input_namespace == namespace, + tbl.c.input_name == name, + ) + ).fetchall() + for r in rows: + entry = dict(r._mapping) + entry["direction"] = "downstream" + results.append(entry) + + return results + + # ── Dataset ownership ── + + def _upsert_dataset_owners( + self, + namespace: str, + name: str, + owners: List[Dict[str, str]], + ): + """Persist all owners from an OwnershipDatasetFacet.""" + now = int(time.time() * 1000) + tbl = OL_TABLES["dataset_ownership"] + + with self._engine.begin() as conn: + for owner in owners: + o_name = owner.get("name", "") + o_type = owner.get("type", "") + if not o_name: + continue + + existing = conn.execute( + select(tbl).where( + tbl.c.dataset_namespace == namespace, + tbl.c.dataset_name == name, + tbl.c.owner_name == o_name, + ) + ).first() + + if existing: + conn.execute( + tbl.update() + .where( + tbl.c.dataset_namespace == namespace, + tbl.c.dataset_name == name, + tbl.c.owner_name == o_name, + ) + .values(owner_type=o_type, updated_at=now) + ) + else: + conn.execute( + tbl.insert().values( + dataset_namespace=namespace, + dataset_name=name, + owner_name=o_name, + owner_type=o_type, + updated_at=now, + ) + ) + + def get_dataset_owners( + self, + namespace: str, + name: str, + ) -> List[Dict[str, str]]: + """Get all owners for a dataset.""" + tbl = OL_TABLES["dataset_ownership"] + with self._engine.connect() as conn: + rows = conn.execute( + select(tbl).where( + tbl.c.dataset_namespace == namespace, + tbl.c.dataset_name == name, + ) + ).fetchall() + return [ + {"name": r._mapping["owner_name"], "type": r._mapping["owner_type"]} + for r in rows + ] + + # ── Execution hierarchy ── + + def get_child_runs( + self, + parent_run_id: str, + limit: int = 50, + ) -> List[Dict[str, Any]]: + """Get all child runs of a given parent run.""" + tbl = OL_TABLES["runs"] + with self._engine.connect() as conn: + rows = conn.execute( + select(tbl) + .where(tbl.c.parent_run_id == parent_run_id) + .order_by(tbl.c.updated_at.desc()) + .limit(limit) + ).fetchall() + return [dict(r._mapping) for r in rows] + + def get_run_tree(self, root_run_id: str) -> List[Dict[str, Any]]: + """Get all runs in an execution tree (sharing the same root).""" + tbl = OL_TABLES["runs"] + with self._engine.connect() as conn: + rows = conn.execute( + select(tbl) + .where( + (tbl.c.root_run_id == root_run_id) | (tbl.c.run_id == root_run_id) + ) + .order_by(tbl.c.started_at.asc().nulls_last()) + ).fetchall() + return [dict(r._mapping) for r in rows] + + # ── Assurance level (computed) ── + + def compute_assurance_level( + self, + namespace: str, + name: str, + ) -> Dict[str, Any]: + """Compute the assurance level for a dataset based on facet presence. + + Levels (cumulative): + - Linked: dataset is in the lineage graph + - Observed: has source evidence (dataSource URI, ETag, hash, etc.) + - Reproducible: inputs+outputs have immutable identifiers or + version snapshots sufficient to reconstruct the run + """ + tbl_ds = OL_TABLES["datasets"] + tbl_ver = OL_TABLES["dataset_versions"] + tbl_edges = OL_TABLES["lineage_edges"] + + result: Dict[str, Any] = { + "namespace": namespace, + "name": name, + "level": "none", + "details": {}, + } + + with self._engine.connect() as conn: + ds_row = conn.execute( + select(tbl_ds).where( + tbl_ds.c.dataset_namespace == namespace, + tbl_ds.c.dataset_name == name, + ) + ).first() + if not ds_row: + return result + + # Level 1: Linked — exists in the graph + has_edges = conn.execute( + select(func.count()) + .select_from(tbl_edges) + .where( + ( + (tbl_edges.c.source_namespace == namespace) + & (tbl_edges.c.source_name == name) + ) + | ( + (tbl_edges.c.target_namespace == namespace) + & (tbl_edges.c.target_name == name) + ) + ) + ).scalar() + + if has_edges and has_edges > 0: + result["level"] = "linked" + result["details"]["edge_count"] = has_edges + else: + return result + + # Level 2: Observed — has source evidence + facets = _safe_parse_json(ds_row._mapping.get("facets_json")) + has_source_evidence = False + evidence_indicators = [] + if facets: + if facets.get("dataSource", {}).get("uri"): + has_source_evidence = True + evidence_indicators.append("dataSource.uri") + if facets.get("storage", {}).get("storageLayer"): + has_source_evidence = True + evidence_indicators.append("storage") + for key in ("dataQualityMetrics", "dataQualityAssertions"): + if key in facets: + has_source_evidence = True + evidence_indicators.append(key) + + if has_source_evidence: + result["level"] = "observed" + result["details"]["evidence"] = evidence_indicators + + # Level 3: Reproducible — has versioned snapshots + version_count = conn.execute( + select(func.count()) + .select_from(tbl_ver) + .where( + tbl_ver.c.dataset_namespace == namespace, + tbl_ver.c.dataset_name == name, + ) + ).scalar() + + if version_count and version_count > 0: + has_schema = ds_row._mapping.get("schema_json") is not None + if has_source_evidence and has_schema: + result["level"] = "reproducible" + result["details"]["version_count"] = version_count + + return result + # ── Cleanup methods ── def delete_dataset(self, namespace: str, name: str): @@ -758,6 +1183,35 @@ def delete_dataset(self, namespace: str, name: str): ) ) + # Clean extended tables + tbl_ver = OL_TABLES["dataset_versions"] + conn.execute( + tbl_ver.delete().where( + (tbl_ver.c.dataset_namespace == namespace) + & (tbl_ver.c.dataset_name == name) + ) + ) + tbl_cl = OL_TABLES["column_lineage"] + conn.execute( + tbl_cl.delete().where( + ( + (tbl_cl.c.dataset_namespace == namespace) + & (tbl_cl.c.dataset_name == name) + ) + | ( + (tbl_cl.c.input_namespace == namespace) + & (tbl_cl.c.input_name == name) + ) + ) + ) + tbl_own = OL_TABLES["dataset_ownership"] + conn.execute( + tbl_own.delete().where( + (tbl_own.c.dataset_namespace == namespace) + & (tbl_own.c.dataset_name == name) + ) + ) + tbl_ds = OL_TABLES["datasets"] conn.execute( tbl_ds.delete().where( @@ -815,6 +1269,9 @@ def delete_job(self, namespace: str, name: str): def purge_all(self): """Delete all data from all OpenLineage tables.""" table_order = [ + "column_lineage", + "dataset_ownership", + "dataset_versions", "run_io", "runs", "lineage_edges", @@ -860,6 +1317,23 @@ def purge_namespace(self, namespace: str): ) ) + # Purge extended tables for this namespace + tbl_ver = OL_TABLES["dataset_versions"] + conn.execute( + tbl_ver.delete().where(tbl_ver.c.dataset_namespace == namespace) + ) + tbl_cl = OL_TABLES["column_lineage"] + conn.execute( + tbl_cl.delete().where( + (tbl_cl.c.dataset_namespace == namespace) + | (tbl_cl.c.input_namespace == namespace) + ) + ) + tbl_own = OL_TABLES["dataset_ownership"] + conn.execute( + tbl_own.delete().where(tbl_own.c.dataset_namespace == namespace) + ) + tbl_ds = OL_TABLES["datasets"] conn.execute(tbl_ds.delete().where(tbl_ds.c.dataset_namespace == namespace)) @@ -970,6 +1444,13 @@ def prune_expired(self, retention_days: int) -> Dict[str, int]: ) deleted["events"] = result.rowcount + # 5. Delete expired dataset versions + tbl_ver = OL_TABLES["dataset_versions"] + result = conn.execute( + tbl_ver.delete().where(tbl_ver.c.created_at < cutoff_ms) + ) + deleted["dataset_versions"] = result.rowcount + total = sum(deleted.values()) if total > 0: logger.info( diff --git a/sdk/python/tests/unit/openlineage/test_processor.py b/sdk/python/tests/unit/openlineage/test_processor.py index c7e3e1a0ac3..cf68cd4168b 100644 --- a/sdk/python/tests/unit/openlineage/test_processor.py +++ b/sdk/python/tests/unit/openlineage/test_processor.py @@ -692,3 +692,388 @@ def test_upstream_from_output(self, processor, store): assert "a" in names assert "j1" in names assert "c" in names + + +# ── Parent hierarchy extraction ── + + +class TestParentHierarchy: + def test_parent_run_id_extracted(self, processor, store): + event = _run_event( + run_id="child-run", + run_facets={ + "parent": { + "run": {"runId": "parent-run"}, + "job": {"namespace": "ns", "name": "parent-job"}, + } + }, + ) + processor.process_event(event) + runs = store.get_runs() + assert len(runs) == 1 + assert runs[0]["parent_run_id"] == "parent-run" + assert runs[0]["root_run_id"] == "parent-run" + + def test_root_run_id_from_nested_root(self, processor, store): + event = _run_event( + run_id="grandchild", + run_facets={ + "parent": { + "run": {"runId": "child-run"}, + "job": {"namespace": "ns", "name": "child-job"}, + "root": { + "run": {"runId": "root-run"}, + "job": {"namespace": "ns", "name": "root-job"}, + }, + } + }, + ) + processor.process_event(event) + runs = store.get_runs() + assert runs[0]["parent_run_id"] == "child-run" + assert runs[0]["root_run_id"] == "root-run" + + def test_no_parent_facet(self, processor, store): + event = _run_event(run_id="standalone") + processor.process_event(event) + runs = store.get_runs() + assert runs[0]["parent_run_id"] is None + assert runs[0]["root_run_id"] is None + + def test_child_runs_query(self, processor, store): + processor.process_event( + _run_event( + job_name="parent-j", + run_id="parent-run", + ) + ) + processor.process_event( + _run_event( + job_name="child-j", + run_id="child-1", + run_facets={ + "parent": { + "run": {"runId": "parent-run"}, + "job": {"namespace": "test-ns", "name": "parent-j"}, + } + }, + ) + ) + processor.process_event( + _run_event( + job_name="child-j2", + run_id="child-2", + run_facets={ + "parent": { + "run": {"runId": "parent-run"}, + "job": {"namespace": "test-ns", "name": "parent-j"}, + } + }, + ) + ) + children = store.get_child_runs("parent-run") + assert len(children) == 2 + child_ids = {c["run_id"] for c in children} + assert "child-1" in child_ids + assert "child-2" in child_ids + + def test_run_tree_query(self, processor, store): + processor.process_event(_run_event(run_id="root")) + for i in range(3): + processor.process_event( + _run_event( + job_name=f"child-{i}", + run_id=f"child-{i}", + run_facets={ + "parent": { + "run": {"runId": "root"}, + "job": {"namespace": "test-ns", "name": "etl-job"}, + "root": { + "run": {"runId": "root"}, + "job": {"namespace": "test-ns", "name": "etl-job"}, + }, + } + }, + ) + ) + tree = store.get_run_tree("root") + assert len(tree) == 4 + run_ids = {r["run_id"] for r in tree} + assert "root" in run_ids + + +# ── Column-level lineage ── + + +class TestColumnLineage: + def test_column_lineage_extracted(self, processor, store): + event = _run_event( + outputs=[ + { + "namespace": "ns", + "name": "output_table", + "facets": { + "columnLineage": { + "fields": { + "full_name": { + "inputFields": [ + { + "namespace": "ns", + "name": "input_table", + "field": "first_name", + "transformations": [ + { + "type": "DIRECT", + "description": "concatenation", + } + ], + }, + { + "namespace": "ns", + "name": "input_table", + "field": "last_name", + }, + ] + } + } + } + }, + } + ], + ) + processor.process_event(event) + cl = store.get_column_lineage("ns", "output_table", direction="upstream") + assert len(cl) == 2 + fields = {(c["input_field"], c["output_field"]) for c in cl} + assert ("first_name", "full_name") in fields + assert ("last_name", "full_name") in fields + + xform = next(c for c in cl if c["input_field"] == "first_name") + assert xform["transformation_type"] == "DIRECT" + + def test_column_lineage_downstream_query(self, processor, store): + event = _run_event( + outputs=[ + { + "namespace": "ns", + "name": "derived", + "facets": { + "columnLineage": { + "fields": { + "score": { + "inputFields": [ + { + "namespace": "ns", + "name": "source", + "field": "raw_score", + } + ] + } + } + } + }, + } + ], + ) + processor.process_event(event) + downstream = store.get_column_lineage("ns", "source", direction="downstream") + assert len(downstream) == 1 + assert downstream[0]["output_field"] == "score" + assert downstream[0]["dataset_name"] == "derived" + + def test_no_column_lineage_when_absent(self, processor, store): + event = _run_event( + outputs=[{"namespace": "ns", "name": "tbl", "facets": {}}], + ) + processor.process_event(event) + cl = store.get_column_lineage("ns", "tbl") + assert len(cl) == 0 + + +# ── Dataset versioning ── + + +class TestDatasetVersioning: + def test_version_created_on_complete(self, processor, store): + event = _run_event( + event_type="COMPLETE", + run_id="ver-run-1", + outputs=[ + { + "namespace": "ns", + "name": "versioned_ds", + "facets": {"schema": {"fields": [{"name": "id", "type": "INT"}]}}, + } + ], + ) + processor.process_event(event) + versions = store.get_dataset_versions("ns", "versioned_ds") + assert len(versions) == 1 + assert versions[0]["version"] == 1 + assert versions[0]["created_by_run_id"] == "ver-run-1" + + datasets = store.get_datasets(namespaces=["ns"]) + ds = next(d for d in datasets if d["dataset_name"] == "versioned_ds") + assert ds["current_version"] == 1 + + def test_multiple_versions(self, processor, store): + for i in range(3): + event = _run_event( + event_type="COMPLETE", + run_id=f"run-{i}", + job_name=f"job-{i}", + outputs=[{"namespace": "ns", "name": "multi_ver", "facets": {}}], + ) + processor.process_event(event) + + versions = store.get_dataset_versions("ns", "multi_ver") + assert len(versions) == 3 + assert versions[0]["version"] == 3 + assert versions[2]["version"] == 1 + + def test_no_version_on_start(self, processor, store): + event = _run_event( + event_type="START", + outputs=[{"namespace": "ns", "name": "ds", "facets": {}}], + ) + processor.process_event(event) + versions = store.get_dataset_versions("ns", "ds") + assert len(versions) == 0 + + def test_get_specific_version(self, processor, store): + for i in range(2): + processor.process_event( + _run_event( + event_type="COMPLETE", + run_id=f"r{i}", + job_name=f"j{i}", + outputs=[{"namespace": "ns", "name": "ds", "facets": {}}], + ) + ) + v1 = store.get_dataset_version("ns", "ds", 1) + assert v1 is not None + assert v1["version"] == 1 + assert store.get_dataset_version("ns", "ds", 99) is None + + +# ── Ownership extraction ── + + +class TestOwnershipExtraction: + def test_ownership_facet_indexed(self, processor, store): + event = _run_event( + outputs=[ + { + "namespace": "ns", + "name": "owned_ds", + "facets": { + "ownership": { + "owners": [ + {"name": "team-ml", "type": "TEAM"}, + {"name": "alice@example.com", "type": "PERSON"}, + ] + } + }, + } + ], + ) + processor.process_event(event) + + datasets = store.get_datasets(namespaces=["ns"]) + ds = next(d for d in datasets if d["dataset_name"] == "owned_ds") + assert ds["owner_name"] == "team-ml" + assert ds["owner_type"] == "TEAM" + + owners = store.get_dataset_owners("ns", "owned_ds") + assert len(owners) == 2 + names = {o["name"] for o in owners} + assert "team-ml" in names + assert "alice@example.com" in names + + def test_no_ownership_when_absent(self, processor, store): + event = _run_event( + outputs=[{"namespace": "ns", "name": "no_owner", "facets": {}}], + ) + processor.process_event(event) + owners = store.get_dataset_owners("ns", "no_owner") + assert len(owners) == 0 + datasets = store.get_datasets(namespaces=["ns"]) + ds = next(d for d in datasets if d["dataset_name"] == "no_owner") + assert ds["owner_name"] is None + + +# ── Lifecycle tracking ── + + +class TestLifecycleTracking: + def test_lifecycle_state_indexed(self, processor, store): + event = _dataset_event( + ds_facets={"lifecycleStateChange": {"lifecycleStateChange": "CREATE"}}, + ) + processor.process_event(event) + datasets = store.get_datasets() + assert datasets[0]["lifecycle_state"] == "CREATE" + + def test_lifecycle_update(self, processor, store): + processor.process_event( + _dataset_event( + ds_facets={"lifecycleStateChange": {"lifecycleStateChange": "CREATE"}}, + ) + ) + processor.process_event( + _dataset_event( + ds_facets={"lifecycleStateChange": {"lifecycleStateChange": "ALTER"}}, + ) + ) + datasets = store.get_datasets() + assert len(datasets) == 1 + assert datasets[0]["lifecycle_state"] == "ALTER" + + +# ── Assurance level computation ── + + +class TestAssuranceLevel: + def test_none_for_unknown_dataset(self, store): + result = store.compute_assurance_level("ns", "nonexistent") + assert result["level"] == "none" + + def test_none_when_no_edges(self, store): + store.upsert_dataset("ns", "isolated") + result = store.compute_assurance_level("ns", "isolated") + assert result["level"] == "none" + + def test_linked_with_edges(self, store): + store.upsert_dataset("ns", "ds1") + store.upsert_lineage_edge("dataset", "ns", "ds1", "job", "ns", "j1") + result = store.compute_assurance_level("ns", "ds1") + assert result["level"] == "linked" + assert result["details"]["edge_count"] >= 1 + + def test_observed_with_source_evidence(self, store): + store.upsert_dataset( + "ns", + "ds2", + facets={"dataSource": {"uri": "s3://bucket/path"}}, + ) + store.upsert_lineage_edge("dataset", "ns", "ds2", "job", "ns", "j1") + result = store.compute_assurance_level("ns", "ds2") + assert result["level"] == "observed" + assert "dataSource.uri" in result["details"]["evidence"] + + def test_reproducible_with_versions(self, store): + store.upsert_dataset( + "ns", + "ds3", + facets={ + "dataSource": {"uri": "s3://bucket/path"}, + "schema": {"fields": [{"name": "id", "type": "INT"}]}, + }, + ) + store.upsert_lineage_edge("dataset", "ns", "ds3", "job", "ns", "j1") + store.create_dataset_version( + "ns", "ds3", run_id="r1", schema_json='{"fields": []}' + ) + result = store.compute_assurance_level("ns", "ds3") + assert result["level"] == "reproducible" + assert result["details"]["version_count"] == 1 diff --git a/sdk/python/tests/unit/openlineage/test_store.py b/sdk/python/tests/unit/openlineage/test_store.py index 66ffb9cd325..4226949c940 100644 --- a/sdk/python/tests/unit/openlineage/test_store.py +++ b/sdk/python/tests/unit/openlineage/test_store.py @@ -686,3 +686,294 @@ def test_retention_stats(self, store): assert stats["jobs"]["count"] == 1 assert stats["datasets"]["count"] == 1 assert "oldest_ms" in stats["events"] + + def test_prune_deletes_old_dataset_versions(self, store): + self._insert_old_and_new(store) + + old_ms = int((time.time() - 60 * 86400) * 1000) + new_ms = int((time.time() - 5 * 86400) * 1000) + tbl_ver = OL_TABLES["dataset_versions"] + store.upsert_dataset("ns", "ds1") + with store.engine.begin() as conn: + conn.execute( + tbl_ver.insert().values( + dataset_namespace="ns", + dataset_name="ds1", + version=1, + created_by_run_id="old-run", + created_at=old_ms, + ) + ) + conn.execute( + tbl_ver.insert().values( + dataset_namespace="ns", + dataset_name="ds1", + version=2, + created_by_run_id="new-run", + created_at=new_ms, + ) + ) + + deleted = store.prune_expired(retention_days=30) + assert deleted["dataset_versions"] == 1 + + versions = store.get_dataset_versions("ns", "ds1") + assert len(versions) == 1 + assert versions[0]["version"] == 2 + + +# ── Dataset versioning ── + + +class TestDatasetVersioning: + def test_create_version(self, store): + store.upsert_dataset("ns", "ds1") + ver = store.create_dataset_version("ns", "ds1", run_id="r1") + assert ver == 1 + + def test_increment_versions(self, store): + store.upsert_dataset("ns", "ds1") + assert store.create_dataset_version("ns", "ds1") == 1 + assert store.create_dataset_version("ns", "ds1") == 2 + assert store.create_dataset_version("ns", "ds1") == 3 + + def test_current_version_updated(self, store): + store.upsert_dataset("ns", "ds1") + store.create_dataset_version("ns", "ds1") + store.create_dataset_version("ns", "ds1") + datasets = store.get_datasets(namespaces=["ns"]) + assert datasets[0]["current_version"] == 2 + + def test_list_versions(self, store): + store.upsert_dataset("ns", "ds1") + for i in range(5): + store.create_dataset_version("ns", "ds1", run_id=f"run-{i}") + versions = store.get_dataset_versions("ns", "ds1", limit=3) + assert len(versions) == 3 + assert versions[0]["version"] == 5 + + def test_get_specific_version(self, store): + store.upsert_dataset("ns", "ds1") + store.create_dataset_version( + "ns", + "ds1", + run_id="r1", + schema_json='{"fields": []}', + facets_json='{"key": "val"}', + ) + v = store.get_dataset_version("ns", "ds1", 1) + assert v is not None + assert v["created_by_run_id"] == "r1" + + def test_get_missing_version(self, store): + store.upsert_dataset("ns", "ds1") + assert store.get_dataset_version("ns", "ds1", 99) is None + + +# ── Column lineage store ── + + +class TestColumnLineageStore: + def test_upsert_and_query(self, store): + store.upsert_column_lineage( + "ns", + "out_ds", + "col_a", + "ns", + "in_ds", + "col_x", + transformation_type="DIRECT", + ) + cl = store.get_column_lineage("ns", "out_ds", direction="upstream") + assert len(cl) == 1 + assert cl[0]["output_field"] == "col_a" + assert cl[0]["input_field"] == "col_x" + assert cl[0]["transformation_type"] == "DIRECT" + + def test_dedup(self, store): + for _ in range(3): + store.upsert_column_lineage( + "ns", + "out_ds", + "col_a", + "ns", + "in_ds", + "col_x", + ) + cl = store.get_column_lineage("ns", "out_ds") + upstream = [c for c in cl if c["direction"] == "upstream"] + assert len(upstream) == 1 + + def test_downstream_query(self, store): + store.upsert_column_lineage( + "ns", + "out_ds", + "col_a", + "ns", + "in_ds", + "col_x", + ) + cl = store.get_column_lineage("ns", "in_ds", direction="downstream") + assert len(cl) == 1 + assert cl[0]["dataset_name"] == "out_ds" + assert cl[0]["direction"] == "downstream" + + +# ── Dataset ownership store ── + + +class TestDatasetOwnershipStore: + def test_upsert_and_query(self, store): + store.upsert_dataset("ns", "owned_ds") + store._upsert_dataset_owners( + "ns", + "owned_ds", + [ + {"name": "alice", "type": "PERSON"}, + {"name": "team-data", "type": "TEAM"}, + ], + ) + owners = store.get_dataset_owners("ns", "owned_ds") + assert len(owners) == 2 + names = {o["name"] for o in owners} + assert "alice" in names + assert "team-data" in names + + def test_upsert_updates_type(self, store): + store._upsert_dataset_owners("ns", "ds1", [{"name": "alice", "type": "PERSON"}]) + store._upsert_dataset_owners("ns", "ds1", [{"name": "alice", "type": "ADMIN"}]) + owners = store.get_dataset_owners("ns", "ds1") + assert len(owners) == 1 + assert owners[0]["type"] == "ADMIN" + + def test_empty_owner_skipped(self, store): + store._upsert_dataset_owners("ns", "ds1", [{"name": "", "type": "PERSON"}]) + owners = store.get_dataset_owners("ns", "ds1") + assert len(owners) == 0 + + +# ── Run hierarchy store ── + + +class TestRunHierarchyStore: + def test_parent_and_root_stored(self, store): + store.upsert_job("ns", "j1", {"facets": {}}) + store.upsert_run( + "child-1", + "ns", + "j1", + "COMPLETE", + parent_run_id="parent-1", + root_run_id="root-1", + ) + runs = store.get_runs() + assert runs[0]["parent_run_id"] == "parent-1" + assert runs[0]["root_run_id"] == "root-1" + + def test_child_runs(self, store): + store.upsert_job("ns", "j1", {"facets": {}}) + store.upsert_run("parent", "ns", "j1", "COMPLETE") + store.upsert_run( + "child-a", + "ns", + "j1", + "COMPLETE", + parent_run_id="parent", + ) + store.upsert_run( + "child-b", + "ns", + "j1", + "COMPLETE", + parent_run_id="parent", + ) + children = store.get_child_runs("parent") + assert len(children) == 2 + + def test_run_tree(self, store): + store.upsert_job("ns", "j1", {"facets": {}}) + store.upsert_run("root", "ns", "j1", "COMPLETE") + store.upsert_run( + "child-1", + "ns", + "j1", + "COMPLETE", + parent_run_id="root", + root_run_id="root", + ) + store.upsert_run( + "grandchild", + "ns", + "j1", + "COMPLETE", + parent_run_id="child-1", + root_run_id="root", + ) + tree = store.get_run_tree("root") + assert len(tree) == 3 + run_ids = {r["run_id"] for r in tree} + assert {"root", "child-1", "grandchild"} == run_ids + + +# ── Purge with extended tables ── + + +class TestPurgeExtendedTables: + def test_purge_all_clears_extended(self, store): + store.upsert_dataset("ns", "ds1") + store.create_dataset_version("ns", "ds1") + store.upsert_column_lineage( + "ns", + "ds1", + "col", + "ns", + "src", + "src_col", + ) + store._upsert_dataset_owners("ns", "ds1", [{"name": "owner", "type": "PERSON"}]) + + store.purge_all() + assert len(store.get_dataset_versions("ns", "ds1")) == 0 + assert len(store.get_column_lineage("ns", "ds1")) == 0 + assert len(store.get_dataset_owners("ns", "ds1")) == 0 + + def test_purge_namespace_clears_extended(self, store): + store.upsert_dataset("ns-a", "ds1") + store.upsert_dataset("ns-b", "ds2") + store.create_dataset_version("ns-a", "ds1") + store.create_dataset_version("ns-b", "ds2") + store.upsert_column_lineage( + "ns-a", + "ds1", + "col", + "ns-a", + "src", + "src_col", + ) + store._upsert_dataset_owners( + "ns-a", "ds1", [{"name": "owner", "type": "PERSON"}] + ) + + store.purge_namespace("ns-a") + assert len(store.get_dataset_versions("ns-a", "ds1")) == 0 + assert len(store.get_column_lineage("ns-a", "ds1")) == 0 + assert len(store.get_dataset_owners("ns-a", "ds1")) == 0 + assert len(store.get_dataset_versions("ns-b", "ds2")) == 1 + + def test_delete_dataset_clears_extended(self, store): + store.upsert_dataset("ns", "ds1") + store.create_dataset_version("ns", "ds1") + store.upsert_column_lineage( + "ns", + "ds1", + "col", + "ns", + "src", + "src_col", + ) + store._upsert_dataset_owners("ns", "ds1", [{"name": "owner", "type": "PERSON"}]) + + store.delete_dataset("ns", "ds1") + assert len(store.get_dataset_versions("ns", "ds1")) == 0 + assert len(store.get_column_lineage("ns", "ds1")) == 0 + assert len(store.get_dataset_owners("ns", "ds1")) == 0 diff --git a/ui/src/components/OpenLineageGraph.tsx b/ui/src/components/OpenLineageGraph.tsx index 9c613ec4ab8..03ad808cdc9 100644 --- a/ui/src/components/OpenLineageGraph.tsx +++ b/ui/src/components/OpenLineageGraph.tsx @@ -931,6 +931,21 @@ const NodeDetailPanel: React.FC<{ {node.source_type} )} + {node.current_version != null && ( + + v{node.current_version} + + )} + {node.lifecycle_state && ( + + {node.lifecycle_state} + + )} + {node.owner && ( + + {node.owner} + + )}
{node.namespace} diff --git a/ui/src/queries/useLoadOpenLineageGraph.ts b/ui/src/queries/useLoadOpenLineageGraph.ts index ec6350e8763..c9209fce304 100644 --- a/ui/src/queries/useLoadOpenLineageGraph.ts +++ b/ui/src/queries/useLoadOpenLineageGraph.ts @@ -16,6 +16,10 @@ export interface OpenLineageNode { description?: string; job_type?: string; source_type?: string; + owner?: string; + owner_type?: string; + lifecycle_state?: string; + current_version?: number; facets?: Record; }