Skip to content

Commit d6bfbb7

Browse files
committed
feat: cover lineage query profile and usage APIs
1 parent 307877c commit d6bfbb7

3 files changed

Lines changed: 182 additions & 0 deletions

File tree

platform-api/src/openmetadata_demo_api/catalog.py

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -84,6 +84,23 @@
8484
"dqtests/TestSuiteResource.java": "child-assets",
8585
"dqtests/TestCaseResource.java": "enrichment",
8686
}
87+
LINEAGE_QUERY_PREFIXES = ("lineage/", "query/", "entityProfiles/", "usage/")
88+
PROFILE_SAMPLE_PREFIXES = ("databases/", "drives/", "searchindex/", "storages/", "topics/")
89+
LINEAGE_QUERY_PREREQUISITES = {
90+
"createDatabaseService",
91+
"createDatabase",
92+
"createDBSchema",
93+
"createTable",
94+
"createDriveService",
95+
"createDirectory",
96+
"createFile",
97+
"createMessagingService",
98+
"createTopic",
99+
"createSearchService",
100+
"createSearchIndex",
101+
"createStorageService",
102+
"createContainer",
103+
}
87104
CORE_SERVICE_PREREQUISITES = {
88105
"createApiService",
89106
"createDashboardService",
@@ -255,6 +272,21 @@ def _is_bulk_drive_operation(operation: Mapping[str, Any]) -> bool:
255272
)
256273

257274

275+
def _is_lineage_query_profile_usage_operation(operation: Mapping[str, Any]) -> bool:
276+
source_file = str(operation["source"]["file"])
277+
text = " ".join(
278+
(str(operation["operation_id"]), str(operation["java_method"]), str(operation["path"]))
279+
).lower()
280+
return (
281+
source_file.startswith(LINEAGE_QUERY_PREFIXES)
282+
or (
283+
source_file.startswith(PROFILE_SAMPLE_PREFIXES)
284+
and any(term in text for term in ("sample", "profile"))
285+
)
286+
or operation["operation_id"] in LINEAGE_QUERY_PREREQUISITES
287+
)
288+
289+
258290
def _phase(operation: Mapping[str, Any]) -> str:
259291
source_file = operation["source"]["file"]
260292
if source_file.startswith("services/"):
@@ -419,6 +451,15 @@ def scenarios() -> dict[str, Scenario]:
419451
or operation["operation_id"] in DATA_QUALITY_PREREQUISITES
420452
),
421453
),
454+
"lineage-query-profile-usage": scenario_from_operations(
455+
"lineage-query-profile-usage",
456+
"Lineage, queries, profiles, usage, sample data, and deterministic replay",
457+
tuple(
458+
operation
459+
for operation in all_operations
460+
if _is_lineage_query_profile_usage_operation(operation)
461+
),
462+
),
422463
}
423464

424465

platform-api/src/openmetadata_demo_api/request_fixtures.py

Lines changed: 62 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,15 +10,45 @@
1010
from typing import Any
1111

1212
DEMO_ID = "00000000-0000-0000-0000-000000000001"
13+
DEMO_ID_2 = "00000000-0000-0000-0000-000000000002"
1314
DEMO_ENTITY_REFERENCE = {
1415
"id": DEMO_ID,
1516
"type": "table",
1617
"name": "demo",
1718
"fullyQualifiedName": "demo.demo.demo.demo",
1819
}
20+
OPENLINEAGE_EVENT = {
21+
"eventTime": "2026-01-01T00:00:00Z",
22+
"producer": "https://open-metadata.org/demo",
23+
"schemaURL": "https://openlineage.io/spec/2-0-2/OpenLineage.json",
24+
"eventType": "COMPLETE",
25+
"run": {"runId": DEMO_ID},
26+
"job": {"namespace": "demo", "name": "daily_orders"},
27+
"inputs": [],
28+
"outputs": [],
29+
}
1930

2031
REQUEST_FIXTURES = MappingProxyType(
2132
{
33+
"AddLineage": {
34+
"edge": {
35+
"fromEntity": {
36+
**DEMO_ENTITY_REFERENCE,
37+
"name": "source",
38+
"fullyQualifiedName": "demo.demo.demo.source",
39+
},
40+
"toEntity": {
41+
**DEMO_ENTITY_REFERENCE,
42+
"id": DEMO_ID_2,
43+
"name": "target",
44+
"fullyQualifiedName": "demo.demo.demo.target",
45+
},
46+
"lineageDetails": {
47+
"source": "Manual",
48+
"description": "Dry-run lineage",
49+
},
50+
}
51+
},
2252
"AIGovernanceBulkTriageRequest": {
2353
"action": "register",
2454
"reason": "Dry-run review",
@@ -85,6 +115,7 @@
85115
}
86116
]
87117
},
118+
"DailyCount": {"count": 1, "date": "2026-01-01"},
88119
"DatabaseProfilerConfig": {"sampleDataCount": 10, "randomizedSample": True},
89120
"DatabaseSchemaProfilerConfig": {
90121
"sampleDataCount": 10,
@@ -100,6 +131,11 @@
100131
"value": "Updated by the dry-run demo",
101132
}
102133
],
134+
"HydrateLineageRequest": {
135+
"entities": [DEMO_ENTITY_REFERENCE],
136+
"fields": "tags,owners",
137+
},
138+
"LineageDetails": {"source": "Manual", "description": "Dry-run lineage"},
103139
"CustomProperty": {
104140
"name": "demoProperty",
105141
"description": "Dry-run custom property",
@@ -137,6 +173,8 @@
137173
}
138174
],
139175
},
176+
"OpenLineageBatchRequest": {"events": [OPENLINEAGE_EVENT]},
177+
"OpenLineageRunEvent": OPENLINEAGE_EVENT,
140178
"PipelineObservability": {
141179
"pipeline": {
142180
**DEMO_ENTITY_REFERENCE,
@@ -220,6 +258,10 @@
220258

221259
PYDANTIC_REQUEST_MODELS = MappingProxyType(
222260
{
261+
"AddLineage": (
262+
"metadata.generated.schema.api.lineage.addLineage",
263+
"AddLineageRequest",
264+
),
223265
"AIGovernanceBulkTriageRequest": (
224266
"metadata.generated.schema.api.ai.aiGovernanceBulkTriageRequest",
225267
"AIGovernanceBulkTriageRequest",
@@ -258,6 +300,10 @@
258300
"metadata.generated.schema.tests.dataQualityReportBatchRequest",
259301
"DataQualityReportBatchRequest",
260302
),
303+
"DailyCount": (
304+
"metadata.generated.schema.type.dailyCount",
305+
"DailyCountOfSomeMeasurement",
306+
),
261307
"DatabaseProfilerConfig": (
262308
"metadata.generated.schema.entity.data.database",
263309
"DatabaseProfilerConfig",
@@ -270,6 +316,14 @@
270316
"metadata.generated.schema.email.emailTemplate",
271317
"EmailTemplate",
272318
),
319+
"HydrateLineageRequest": (
320+
"metadata.generated.schema.api.lineage.hydrateLineageRequest",
321+
"HydrateLineageRequest",
322+
),
323+
"LineageDetails": (
324+
"metadata.generated.schema.type.entityLineage",
325+
"LineageDetails",
326+
),
273327
"EntityReference": (
274328
"metadata.generated.schema.type.entityReference",
275329
"EntityReference",
@@ -290,6 +344,14 @@
290344
"metadata.generated.schema.entity.services.ingestionPipelines.operationMetrics",
291345
"OperationMetricsBatch",
292346
),
347+
"OpenLineageBatchRequest": (
348+
"metadata.generated.schema.api.lineage.openlineage.openLineageBatchRequest",
349+
"OpenLineageBatchRequest",
350+
),
351+
"OpenLineageRunEvent": (
352+
"metadata.generated.schema.api.lineage.openlineage.openLineageRunEvent",
353+
"OpenLineageRunEvent",
354+
),
293355
"PipelineObservability": (
294356
"metadata.generated.schema.type.pipelineObservability",
295357
"PipelineObservability",

platform-api/tests/test_runtime.py

Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -936,3 +936,82 @@ def test_data_quality_requests_are_validated_and_placeholder_free() -> None:
936936
"CreateTestCaseResult",
937937
"CreateTestCaseResolutionStatus",
938938
} <= models
939+
940+
941+
def test_lineage_query_profile_usage_scenario_covers_source_and_dependent_routes() -> None:
942+
prefixes = ("lineage/", "query/", "entityProfiles/", "usage/")
943+
prerequisites = {
944+
"createDatabaseService",
945+
"createDatabase",
946+
"createDBSchema",
947+
"createTable",
948+
"createDriveService",
949+
"createDirectory",
950+
"createFile",
951+
"createMessagingService",
952+
"createTopic",
953+
"createSearchService",
954+
"createSearchIndex",
955+
"createStorageService",
956+
"createContainer",
957+
}
958+
959+
def selected(operation) -> bool:
960+
source_file = str(operation["source"]["file"])
961+
text = " ".join(
962+
(str(operation["operation_id"]), str(operation["java_method"]), str(operation["path"]))
963+
).lower()
964+
return source_file.startswith(prefixes) or (
965+
source_file.startswith(
966+
("databases/", "drives/", "searchindex/", "storages/", "topics/")
967+
)
968+
and any(term in text for term in ("sample", "profile"))
969+
)
970+
971+
expected = {operation["operation_id"] for operation in load_operations() if selected(operation)}
972+
scenario = scenarios()["lineage-query-profile-usage"]
973+
974+
assert {step.operation_id for step in scenario.operation_steps} == expected | prerequisites
975+
assert len(scenario.run()) == len(expected | prerequisites)
976+
assert {
977+
"addLineageEdge",
978+
"addLineageEdgeByName",
979+
"postOpenLineageEvent",
980+
"createQuery",
981+
"addProfileDataById",
982+
"reportEntityUsageWithID",
983+
"addSampleData",
984+
"addDataProfiler",
985+
} <= expected
986+
ordered = [step.operation_id for step in scenario.ordered_steps()]
987+
for dependent in (
988+
"addLineageEdge",
989+
"addLineageEdgeByName",
990+
"reportEntityUsageWithID",
991+
"addSampleData",
992+
"addDataProfiler",
993+
):
994+
assert ordered.index("createTable") < ordered.index(dependent)
995+
assert scenario.run() == scenario.run()
996+
997+
998+
def test_lineage_query_profile_usage_requests_are_validated_and_placeholder_free() -> None:
999+
scenario = scenarios()["lineage-query-profile-usage"]
1000+
1001+
assert not {
1002+
str(value["_model"])
1003+
for step in scenario.operation_steps
1004+
for value in _nested_mappings(step.request)
1005+
if "_model" in value
1006+
}
1007+
models = {
1008+
str(asset["model"]).rsplit(".", 1)[-1]
1009+
for asset in load_assets()
1010+
if "lineage-query-profile-usage" in asset["scenarios"]
1011+
}
1012+
assert {
1013+
"CreateEntityProfileRequest",
1014+
"CreateQueryRequest",
1015+
"CreateQueryCostRecordRequest",
1016+
"CreateTableProfileRequest",
1017+
} <= models

0 commit comments

Comments
 (0)