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 && (
+