From e1c2073d99607a5dc13312ccd0c6af19b0e3e527 Mon Sep 17 00:00:00 2001 From: "a.metwalli" Date: Fri, 17 Jul 2026 23:28:51 +0200 Subject: [PATCH 1/9] feat(contracts): version domain and provider contracts --- docs/architecture/provider-contracts.md | 20 ++++++ docs/schema/domain.schema.json | 52 +++++++++++++++ domains/README.md | 7 +- domains/sales/domain.yaml | 20 +----- domains/supply_chain/domain.yaml | 11 +--- .../environments/aws-poc/contracts.tf | 1 + .../environments/azure-poc/contracts.tf | 1 + .../terraform/environments/local/contracts.tf | 1 + scripts/test/check-contracts.sh | 15 +++++ tools/olf/olf/contracts.py | 22 ++++++- tools/olf/olf/descriptors.py | 64 ++++++++++++++++++ .../fixtures/aws-provider-contracts.json | 2 +- .../fixtures/local-provider-contracts.json | 1 + tools/olf/tests/test_contracts.py | 9 +++ tools/olf/tests/test_descriptors.py | 65 +++++++++++++++++++ 15 files changed, 261 insertions(+), 30 deletions(-) create mode 100644 docs/schema/domain.schema.json create mode 100644 tools/olf/olf/descriptors.py create mode 100644 tools/olf/tests/test_descriptors.py diff --git a/docs/architecture/provider-contracts.md b/docs/architecture/provider-contracts.md index 2cb96ef..30152c8 100644 --- a/docs/architecture/provider-contracts.md +++ b/docs/architecture/provider-contracts.md @@ -9,6 +9,10 @@ state, ingress/TLS, or production hardening services. ## Contract Source Of Truth +Provider contract exports carry `schema_version: "1.0.0"`. Consumers reject +unknown versions rather than guessing field meanings. This version is +independent of deployment and component versions. + Terraform is the source of truth for provider contracts. The local, Azure, and AWS platform roots normalize explicit contract objects in their `contracts.tf` files and validate them with Terraform `check` blocks. Runtime scripts can read @@ -63,6 +67,22 @@ Product-owned runtime assets use logical aliases. Local and Azure resolve AWS resolves them to S3 buckets. Local and Azure resolve `iceberg_catalog` to Polaris; AWS resolves it to Glue. +Domain descriptors use `apiVersion: openlakeforge.io/v1alpha1` and +`kind: Domain`; the machine-readable schema is +[`docs/schema/domain.schema.json`](../schema/domain.schema.json). Descriptors +contain logical product and table names only. Provider contracts derive the +physical catalog/database/schema FQNs at runtime, so changing catalog adapters +does not require editing business metadata. + +### Compatibility and migration + +The `v1alpha1` descriptor and provider contract versions are strict. Migrate a +legacy descriptor by adding the version envelope, removing +`silver_tables.schema`, `gold_tables.schema`, and physical asset FQNs, then +validating against the schema before deployment. A future incompatible shape +must publish a new API/version and migration guide; deployments fail closed +when either version is unknown. + ## Catalog Contract The catalog contract describes an Iceberg catalog implementation. The local diff --git a/docs/schema/domain.schema.json b/docs/schema/domain.schema.json new file mode 100644 index 0000000..80bf098 --- /dev/null +++ b/docs/schema/domain.schema.json @@ -0,0 +1,52 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://openlakeforge.io/schemas/domain/v1alpha1", + "title": "OpenLakeForge Domain", + "type": "object", + "additionalProperties": true, + "required": ["apiVersion", "kind", "name", "displayName", "description", "status", "data_products"], + "properties": { + "apiVersion": {"const": "openlakeforge.io/v1alpha1"}, + "kind": {"const": "Domain"}, + "name": {"type": "string", "pattern": "^[a-z][a-z0-9_]*$"}, + "displayName": {"type": "string", "minLength": 1}, + "description": {"type": "string"}, + "status": {"type": "string"}, + "data_products": { + "type": "array", + "items": { + "type": "object", + "required": ["id", "name", "displayName", "description", "status"], + "properties": { + "id": {"type": "string"}, + "name": {"type": "string"}, + "displayName": {"type": "string"}, + "description": {"type": "string"}, + "status": {"type": "string"}, + "silver_tables": {"$ref": "#/$defs/tableGroup"}, + "gold_tables": {"$ref": "#/$defs/tableGroup"}, + "assets": { + "type": "array", + "items": { + "type": "object", + "not": { + "anyOf": [ + {"required": ["fqn"]}, + {"required": ["fullyQualifiedName"]} + ] + } + } + } + } + } + } + }, + "$defs": { + "tableGroup": { + "type": "object", + "required": ["tables"], + "properties": {"tables": {"type": "array"}}, + "not": {"required": ["schema"]} + } + } +} diff --git a/domains/README.md b/domains/README.md index 37fb4f9..40ef944 100644 --- a/domains/README.md +++ b/domains/README.md @@ -7,7 +7,12 @@ product has side-by-side assets under domain capability folders: raw examples, dlt loaders, Floe contracts, dbt projects, Dagster modules, Superset reports, tests, and documentation. -`domains//domain.yaml` is the single source for domain, data-product, +`domains//domain.yaml` is a versioned (`openlakeforge.io/v1alpha1`, +`kind: Domain`) single source for domain, data-product, Bronze, Silver, Gold, and OpenMetadata table metadata. +Keep this metadata provider-neutral: catalog/database/schema identities are +derived from the environment provider contract. Validate descriptors against +`docs/schema/domain.schema.json` before deployment. + The current seed domains are `sales` and `supply_chain`. diff --git a/domains/sales/domain.yaml b/domains/sales/domain.yaml index 5a1ae5c..e5bd740 100644 --- a/domains/sales/domain.yaml +++ b/domains/sales/domain.yaml @@ -1,3 +1,5 @@ +apiVersion: openlakeforge.io/v1alpha1 +kind: Domain name: sales displayName: Sales domainType: Source-aligned @@ -38,7 +40,6 @@ data_products: path: s3://lakehouse-bronze/sales/order_revenue/promotions description: Raw CSV promotion dimension. silver_tables: - schema: polaris.lakehouse_dev.sales_order_revenue_silver tables: - name: orders description: Validated sales order headers. @@ -51,7 +52,6 @@ data_products: - name: promotions description: Validated promotion dimension. gold_tables: - schema: polaris.lakehouse_dev.sales_order_revenue_gold tables: - name: mart_order_revenue_by_day description: Daily order revenue, units, and discount amount by region. @@ -59,13 +59,6 @@ data_products: description: Net revenue and discount amount by sales channel and promotion type. - name: mart_order_revenue_margin_by_product description: Product revenue, cost, and gross margin for fulfilled order lines. - assets: - - type: table - fqn: polaris.lakehouse_dev.sales_order_revenue_gold.mart_order_revenue_by_day - - type: table - fqn: polaris.lakehouse_dev.sales_order_revenue_gold.mart_order_revenue_by_channel - - type: table - fqn: polaris.lakehouse_dev.sales_order_revenue_gold.mart_order_revenue_margin_by_product - id: customer_health name: sales_customer_health displayName: Sales Customer Health @@ -86,7 +79,6 @@ data_products: path: s3://lakehouse-bronze/sales/customer_health/nps_responses description: Raw CSV NPS responses. silver_tables: - schema: polaris.lakehouse_dev.sales_customer_health_silver tables: - name: accounts description: Validated account dimension. @@ -97,7 +89,6 @@ data_products: - name: nps_responses description: Validated NPS response facts. gold_tables: - schema: polaris.lakehouse_dev.sales_customer_health_gold tables: - name: mart_customer_health_score description: Account-level health score using subscription, support, and NPS signals. @@ -105,10 +96,3 @@ data_products: description: Churn risk and ARR exposure grouped by account segment and region. - name: mart_support_sla_by_customer description: Support volume, resolution hours, and SLA rate by customer and priority. - assets: - - type: table - fqn: polaris.lakehouse_dev.sales_customer_health_gold.mart_customer_health_score - - type: table - fqn: polaris.lakehouse_dev.sales_customer_health_gold.mart_churn_risk_by_segment - - type: table - fqn: polaris.lakehouse_dev.sales_customer_health_gold.mart_support_sla_by_customer diff --git a/domains/supply_chain/domain.yaml b/domains/supply_chain/domain.yaml index c87ecc4..4af7d9c 100644 --- a/domains/supply_chain/domain.yaml +++ b/domains/supply_chain/domain.yaml @@ -1,3 +1,5 @@ +apiVersion: openlakeforge.io/v1alpha1 +kind: Domain name: supply_chain displayName: Supply Chain domainType: Source-aligned @@ -41,7 +43,6 @@ data_products: path: s3://lakehouse-bronze/supply_chain/inventory_reliability/stockout_events description: Raw CSV stockout events. silver_tables: - schema: polaris.lakehouse_dev.supply_chain_inventory_reliability_silver tables: - name: warehouses description: Validated warehouse dimension. @@ -56,7 +57,6 @@ data_products: - name: stockout_events description: Validated stockout event facts. gold_tables: - schema: polaris.lakehouse_dev.supply_chain_inventory_reliability_gold tables: - name: mart_inventory_position description: Latest inventory position and reorder status by product and warehouse. @@ -64,10 +64,3 @@ data_products: description: Supplier delivery lateness, fill rate, and on-time performance. - name: mart_stockout_risk description: Stockout risk by product and warehouse using current inventory and stockout history. - assets: - - type: table - fqn: polaris.lakehouse_dev.supply_chain_inventory_reliability_gold.mart_inventory_position - - type: table - fqn: polaris.lakehouse_dev.supply_chain_inventory_reliability_gold.mart_supplier_delivery_reliability - - type: table - fqn: polaris.lakehouse_dev.supply_chain_inventory_reliability_gold.mart_stockout_risk diff --git a/infra/terraform/environments/aws-poc/contracts.tf b/infra/terraform/environments/aws-poc/contracts.tf index f9f16e0..b37bfa8 100644 --- a/infra/terraform/environments/aws-poc/contracts.tf +++ b/infra/terraform/environments/aws-poc/contracts.tf @@ -291,6 +291,7 @@ locals { } provider_contracts = { + schema_version = "1.0.0" foundation = local.foundation_contract kubernetes_platform = local.kubernetes_platform_contract cluster = local.kubernetes_platform_contract diff --git a/infra/terraform/environments/azure-poc/contracts.tf b/infra/terraform/environments/azure-poc/contracts.tf index d307539..3e47ced 100644 --- a/infra/terraform/environments/azure-poc/contracts.tf +++ b/infra/terraform/environments/azure-poc/contracts.tf @@ -300,6 +300,7 @@ locals { } provider_contracts = { + schema_version = "1.0.0" foundation = local.foundation_contract kubernetes_platform = local.kubernetes_platform_contract cluster = local.kubernetes_platform_contract diff --git a/infra/terraform/environments/local/contracts.tf b/infra/terraform/environments/local/contracts.tf index 5b0d97b..acff0dd 100644 --- a/infra/terraform/environments/local/contracts.tf +++ b/infra/terraform/environments/local/contracts.tf @@ -284,6 +284,7 @@ locals { } provider_contracts = { + schema_version = "1.0.0" foundation = local.foundation_contract kubernetes_platform = local.kubernetes_platform_contract cluster = local.kubernetes_platform_contract diff --git a/scripts/test/check-contracts.sh b/scripts/test/check-contracts.sh index 237b3f9..ef070f8 100644 --- a/scripts/test/check-contracts.sh +++ b/scripts/test/check-contracts.sh @@ -35,6 +35,21 @@ azure_main_text = azure_main_tf.read_text() aws_main_tf = Path("infra/terraform/environments/aws-poc/main.tf") aws_main_text = aws_main_tf.read_text() +for descriptor_path in sorted(Path("domains").glob("*/domain.yaml")): + descriptor_text = descriptor_path.read_text() + if "apiVersion: openlakeforge.io/v1alpha1" not in descriptor_text or "kind: Domain" not in descriptor_text: + errors.append(f"{descriptor_path}: unsupported or missing domain descriptor version") + if "polaris.lakehouse_dev" in descriptor_text or "fqn:" in descriptor_text: + errors.append(f"{descriptor_path}: physical catalog identities must come from provider contracts") + +for contracts_path, contracts_body in [ + (contracts_tf, text), + (azure_contracts_tf, azure_text), + (aws_contracts_tf, aws_text), +]: + if 'schema_version = "1.0.0"' not in contracts_body: + errors.append(f"{contracts_path}: provider contracts must declare schema_version 1.0.0") + required_contracts = [ "foundation_contract", "kubernetes_platform_contract", diff --git a/tools/olf/olf/contracts.py b/tools/olf/olf/contracts.py index d51046f..e96f2a4 100644 --- a/tools/olf/olf/contracts.py +++ b/tools/olf/olf/contracts.py @@ -25,6 +25,12 @@ from collections.abc import Mapping from typing import Any +PROVIDER_CONTRACT_SCHEMA_VERSION = "1.0.0" + + +class ProviderContractError(ValueError): + """Raised when Terraform returns an unsupported provider contract version.""" + # Kept as literal compact-JSON strings for byte parity with the previous bash # defaults. _DEFAULT_CATALOG_NAMESPACES_JSON = ( @@ -67,7 +73,15 @@ def load_provider_contracts(terraform_dir: str) -> dict[str, Any] | None: contracts = json.loads(result.stdout) except json.JSONDecodeError: return None - return contracts if isinstance(contracts, dict) else None + if not isinstance(contracts, dict): + return None + schema_version = contracts.get("schema_version") + if schema_version != PROVIDER_CONTRACT_SCHEMA_VERSION: + raise ProviderContractError( + f"provider_contracts.schema_version {schema_version!r} is unsupported; " + f"expected {PROVIDER_CONTRACT_SCHEMA_VERSION!r}" + ) + return contracts class _Env: @@ -350,6 +364,12 @@ def build_contract_env( env = _Env(base) _apply_default_contract_env(env, base) if contracts is not None: + schema_version = contracts.get("schema_version") + if schema_version != PROVIDER_CONTRACT_SCHEMA_VERSION: + raise ProviderContractError( + f"provider_contracts.schema_version {schema_version!r} is unsupported; " + f"expected {PROVIDER_CONTRACT_SCHEMA_VERSION!r}" + ) _apply_provider_contracts(env, contracts) _apply_default_contract_env(env, base) diff --git a/tools/olf/olf/descriptors.py b/tools/olf/olf/descriptors.py new file mode 100644 index 0000000..58240f0 --- /dev/null +++ b/tools/olf/olf/descriptors.py @@ -0,0 +1,64 @@ +"""Versioned domain descriptor validation and migration helpers.""" + +from __future__ import annotations + +from collections.abc import Mapping +from pathlib import Path +from typing import Any + +import yaml + +DOMAIN_API_VERSION = "openlakeforge.io/v1alpha1" +DOMAIN_KIND = "Domain" + + +class DomainDescriptorError(ValueError): + """Raised when a domain descriptor is missing or uses an unsupported version.""" + + +def validate_domain_descriptor(document: Mapping[str, Any], *, source: str = "domain.yaml") -> None: + """Validate the stable envelope and provider-neutral product metadata.""" + if document.get("apiVersion") != DOMAIN_API_VERSION: + raise DomainDescriptorError( + f"{source}: unsupported apiVersion {document.get('apiVersion')!r}; expected {DOMAIN_API_VERSION!r}" + ) + if document.get("kind") != DOMAIN_KIND: + raise DomainDescriptorError(f"{source}: kind must be {DOMAIN_KIND!r}") + for field in ("name", "displayName", "description", "status", "data_products"): + if field not in document: + raise DomainDescriptorError(f"{source}: missing required field {field!r}") + if not isinstance(document["data_products"], list): + raise DomainDescriptorError(f"{source}: data_products must be an array") + for index, product in enumerate(document["data_products"]): + if not isinstance(product, Mapping): + raise DomainDescriptorError(f"{source}: data_products[{index}] must be an object") + for field in ("id", "name", "displayName", "description", "status"): + if field not in product: + raise DomainDescriptorError(f"{source}: data_products[{index}] missing {field!r}") + for group in ("silver_tables", "gold_tables"): + spec = product.get(group) or {} + if "schema" in spec: + raise DomainDescriptorError(f"{source}: {group}.schema must be derived from provider contracts") + for asset_index, asset in enumerate(product.get("assets") or []): + if isinstance(asset, str): + raise DomainDescriptorError( + f"{source}: data_products[{index}].assets[{asset_index}] must be a logical asset object" + ) + if not isinstance(asset, Mapping): + raise DomainDescriptorError( + f"{source}: data_products[{index}].assets[{asset_index}] must be an object" + ) + if "fqn" in asset or "fullyQualifiedName" in asset: + raise DomainDescriptorError( + f"{source}: data_products[{index}].assets[{asset_index}] must not contain physical FQNs" + ) + + +def load_domain_descriptor(path: str | Path) -> dict[str, Any]: + source = str(path) + with Path(path).open(encoding="utf-8") as handle: + document = yaml.safe_load(handle) + if not isinstance(document, dict): + raise DomainDescriptorError(f"{source}: descriptor must contain a YAML object") + validate_domain_descriptor(document, source=source) + return document diff --git a/tools/olf/tests/fixtures/aws-provider-contracts.json b/tools/olf/tests/fixtures/aws-provider-contracts.json index f4a4f13..bfd9a6c 100644 --- a/tools/olf/tests/fixtures/aws-provider-contracts.json +++ b/tools/olf/tests/fixtures/aws-provider-contracts.json @@ -1 +1 @@ -{"storage":{"logical_name":"lakehouse_storage","provider":"aws","implementation":"storage.aws_s3","protocol":"s3","bronze_bucket_name":"openlakeforge-poc-bronze","silver_bucket_name":"openlakeforge-poc-silver","gold_bucket_name":"openlakeforge-poc-gold","region":"eu-west-1","path_style_access":false,"ssl_mode":"required"},"catalog":{"logical_name":"iceberg_catalog","implementation":"catalog.aws_glue","catalog_type":"glue","catalog_provider":"aws-glue","catalog_name":"lakehouse_dev","runtime_profile":"aws-glue-rest","glue_region":"eu-west-1","glue_catalog_id":"123456789012","glue_rest_uri":"https://glue.eu-west-1.amazonaws.com/iceberg","glue_rest_warehouse":"123456789012","glue_warehouse_prefix":"warehouse/iceberg","catalog_namespace_model":"product-layer","catalog_database_fqn":"aws_glue.lakehouse_dev","silver_schema_fqns":{"sales_order_revenue":"aws_glue.lakehouse_dev.sales_order_revenue_silver"},"gold_schema_fqns":{"sales_order_revenue":"aws_glue.lakehouse_dev.sales_order_revenue_gold"}},"artifact_bucket":{"bucket_name":"openlakeforge-poc-ops","artifact_base_uri":"s3://openlakeforge-poc-ops","access_mode":"remote","base_uri":"s3://openlakeforge-poc-ops/floe/manifests","floe_report_base_uri":"s3://openlakeforge-poc-ops/floe/reports","log_base_uri":"s3://openlakeforge-poc-ops/logs","run_artifact_base_uri":"s3://openlakeforge-poc-ops/run-artifacts","local_upload_access_mode":"direct"},"kubernetes_platform":{"namespace":"lakehouse"},"query":{"catalog_name":"iceberg","endpoint":"http://trino:8080"}} +{"schema_version":"1.0.0","storage":{"logical_name":"lakehouse_storage","provider":"aws","implementation":"storage.aws_s3","protocol":"s3","bronze_bucket_name":"openlakeforge-poc-bronze","silver_bucket_name":"openlakeforge-poc-silver","gold_bucket_name":"openlakeforge-poc-gold","region":"eu-west-1","path_style_access":false,"ssl_mode":"required"},"catalog":{"logical_name":"iceberg_catalog","implementation":"catalog.aws_glue","catalog_type":"glue","catalog_provider":"aws-glue","catalog_name":"lakehouse_dev","runtime_profile":"aws-glue-rest","glue_region":"eu-west-1","glue_catalog_id":"123456789012","glue_rest_uri":"https://glue.eu-west-1.amazonaws.com/iceberg","glue_rest_warehouse":"123456789012","glue_warehouse_prefix":"warehouse/iceberg","catalog_namespace_model":"product-layer","catalog_database_fqn":"aws_glue.lakehouse_dev","silver_schema_fqns":{"sales_order_revenue":"aws_glue.lakehouse_dev.sales_order_revenue_silver"},"gold_schema_fqns":{"sales_order_revenue":"aws_glue.lakehouse_dev.sales_order_revenue_gold"}},"artifact_bucket":{"bucket_name":"openlakeforge-poc-ops","artifact_base_uri":"s3://openlakeforge-poc-ops","access_mode":"remote","base_uri":"s3://openlakeforge-poc-ops/floe/manifests","floe_report_base_uri":"s3://openlakeforge-poc-ops/floe/reports","log_base_uri":"s3://openlakeforge-poc-ops/logs","run_artifact_base_uri":"s3://openlakeforge-poc-ops/run-artifacts","local_upload_access_mode":"direct"},"kubernetes_platform":{"namespace":"lakehouse"},"query":{"catalog_name":"iceberg","endpoint":"http://trino:8080"}} diff --git a/tools/olf/tests/fixtures/local-provider-contracts.json b/tools/olf/tests/fixtures/local-provider-contracts.json index 6715b00..ba00a00 100644 --- a/tools/olf/tests/fixtures/local-provider-contracts.json +++ b/tools/olf/tests/fixtures/local-provider-contracts.json @@ -1,4 +1,5 @@ { + "schema_version": "1.0.0", "storage": { "logical_name": "lakehouse_storage", "provider": "local", diff --git a/tools/olf/tests/test_contracts.py b/tools/olf/tests/test_contracts.py index 9acea03..7df7163 100644 --- a/tools/olf/tests/test_contracts.py +++ b/tools/olf/tests/test_contracts.py @@ -46,6 +46,15 @@ def test_defaults_without_contracts_match_local_profile() -> None: assert unsets == [] +def test_unsupported_provider_contract_version_is_rejected() -> None: + from olf.contracts import ProviderContractError + + contracts = load_fixture("local-provider-contracts.json") + contracts["schema_version"] = "9.0.0" + with pytest.raises(ProviderContractError, match="unsupported"): + build_contract_env({}, contracts) + + def test_local_contracts_apply_seaweedfs_values() -> None: exports, unsets = build_contract_env({}, load_fixture("local-provider-contracts.json")) assert exports["OPENLAKEFORGE_STORAGE_PROVIDER"] == "local" diff --git a/tools/olf/tests/test_descriptors.py b/tools/olf/tests/test_descriptors.py new file mode 100644 index 0000000..dbda16b --- /dev/null +++ b/tools/olf/tests/test_descriptors.py @@ -0,0 +1,65 @@ +from pathlib import Path + +import pytest + +from olf.descriptors import DomainDescriptorError, load_domain_descriptor, validate_domain_descriptor + +ROOT = Path(__file__).parents[3] + + +@pytest.mark.parametrize("path", sorted((ROOT / "domains").glob("*/domain.yaml"))) +def test_seed_domain_descriptors_are_versioned_and_provider_neutral(path: Path) -> None: + descriptor = load_domain_descriptor(path) + assert descriptor["apiVersion"] == "openlakeforge.io/v1alpha1" + assert descriptor["kind"] == "Domain" + + +def test_domain_descriptor_rejects_unsupported_version() -> None: + with pytest.raises(DomainDescriptorError, match="unsupported apiVersion"): + validate_domain_descriptor({"apiVersion": "openlakeforge.io/v2", "kind": "Domain"}) + + +def test_domain_descriptor_rejects_physical_catalog_identity() -> None: + descriptor = { + "apiVersion": "openlakeforge.io/v1alpha1", + "kind": "Domain", + "name": "sales", + "displayName": "Sales", + "description": "Sales", + "status": "planned", + "data_products": [ + { + "id": "orders", + "name": "sales_orders", + "displayName": "Orders", + "description": "Orders", + "status": "planned", + "silver_tables": {"schema": "polaris.lakehouse_dev.sales_orders_silver", "tables": []}, + } + ], + } + with pytest.raises(DomainDescriptorError, match="derived from provider contracts"): + validate_domain_descriptor(descriptor) + + +def test_domain_descriptor_rejects_legacy_string_asset() -> None: + descriptor = { + "apiVersion": "openlakeforge.io/v1alpha1", + "kind": "Domain", + "name": "sales", + "displayName": "Sales", + "description": "Sales", + "status": "planned", + "data_products": [ + { + "id": "orders", + "name": "sales_orders", + "displayName": "Orders", + "description": "Orders", + "status": "planned", + "assets": ["polaris.lakehouse_dev.sales_orders_gold.orders"], + } + ], + } + with pytest.raises(DomainDescriptorError, match="logical asset object"): + validate_domain_descriptor(descriptor) From d4ca48d5f0c89fabc7f85723e306a5a1188f947c Mon Sep 17 00:00:00 2001 From: "a.metwalli" Date: Sat, 18 Jul 2026 00:02:11 +0200 Subject: [PATCH 2/9] fix: validate descriptors before metadata deployment --- tools/olf/olf/openmetadata.py | 37 +++++++++++++++++++------ tools/olf/tests/test_openmetadata.py | 41 ++++++++++++++++++++++++++-- 2 files changed, 67 insertions(+), 11 deletions(-) diff --git a/tools/olf/olf/openmetadata.py b/tools/olf/olf/openmetadata.py index 1150c43..aebbc43 100644 --- a/tools/olf/olf/openmetadata.py +++ b/tools/olf/olf/openmetadata.py @@ -18,7 +18,7 @@ from dataclasses import dataclass, field from pathlib import Path -import yaml +from olf.descriptors import load_domain_descriptor class OpenMetadataError(RuntimeError): @@ -265,12 +265,32 @@ def provider_asset_fqn(self, product: dict, fqn): return f"{schema_fqn}.{table_name}" return fqn + def logical_asset_fqn(self, product: dict, asset: dict) -> str: + """Resolve a provider-neutral table name through the contract schemas.""" + name = asset.get("name") + if not name: + raise OpenMetadataError(f"OpenMetadata table asset is missing 'name' or 'fqn': {asset!r}") + matches = [ + f"{schema_fqn}.{name}" + for schema_fqn, table in self.product_table_specs(product) + if table.get("name") == name + ] + if not matches: + raise OpenMetadataError( + f"OpenMetadata logical table asset '{name}' is not declared in the product table contract." + ) + if len(matches) > 1: + raise OpenMetadataError(f"OpenMetadata logical table asset '{name}' is ambiguous: {matches!r}") + return matches[0] + def asset_with_provider_fqn(self, product: dict, asset): if isinstance(asset, str): return self.provider_asset_fqn(product, asset) if isinstance(asset, dict): rewritten = dict(asset) fqn = rewritten.get("fqn") or rewritten.get("fullyQualifiedName") + if not fqn and rewritten.get("name"): + fqn = self.logical_asset_fqn(product, rewritten) if fqn: rewritten["fqn"] = self.provider_asset_fqn(product, fqn) rewritten.pop("fullyQualifiedName", None) @@ -293,12 +313,16 @@ def product_asset_entries(self, product: dict): fqn = asset.get("fqn") or asset.get("fullyQualifiedName") else: fqn = None - fqn = self.provider_asset_fqn(product, fqn) + resolved = self.asset_with_provider_fqn(product, asset) + if isinstance(resolved, dict): + fqn = resolved.get("fqn") + else: + fqn = self.provider_asset_fqn(product, fqn) if fqn and fqn in seen: continue if fqn: seen.add(fqn) - yield self.asset_with_provider_fqn(product, asset) + yield resolved def storage_bucket_specs(self): specs = [ @@ -542,7 +566,7 @@ def deploy(self) -> None: self.wait_for_openmetadata() self.login() - domain_specs = [(path, _load_yaml(path)) for path in self.domain_files()] + domain_specs = [(path, load_domain_descriptor(path)) for path in self.domain_files()] if not domain_specs: raise OpenMetadataError( f"No OpenMetadata domain metadata files found under {self.config.metadata_root}//domain.yaml" @@ -607,8 +631,3 @@ def deploy(self) -> None: print(f"WARN: {guidance}", file=sys.stderr) else: raise OpenMetadataError(guidance) - - -def _load_yaml(path: Path) -> dict: - with path.open("r", encoding="utf-8") as handle: - return yaml.safe_load(handle) or {} diff --git a/tools/olf/tests/test_openmetadata.py b/tools/olf/tests/test_openmetadata.py index 306f648..f898e7f 100644 --- a/tools/olf/tests/test_openmetadata.py +++ b/tools/olf/tests/test_openmetadata.py @@ -149,14 +149,51 @@ def test_product_assets_use_provider_schema_fqns_and_dedup() -> None: ] +def test_logical_asset_name_resolves_through_provider_contract() -> None: + cfg = om.OpenMetadataConfig.from_environment( + { + "OPENLAKEFORGE_CATALOG_GOLD_SCHEMA_FQNS_JSON": + '{"sales_order_revenue": "aws_glue.lakehouse_dev.sales_order_revenue_gold"}', + }, + base_url="http://x", + admin_email="a", + admin_password="p", + metadata_root="domains", + metadata_source_dir="", + allow_missing_assets=False, + catalog_service="aws_glue", + catalog_database="lakehouse_dev", + cleanup_legacy_default_database=False, + ) + deployer = om.OpenMetadataDeployer(cfg, om.OpenMetadataClient(cfg.base_url)) + product = { + "name": "sales_order_revenue", + "gold_tables": {"tables": [{"name": "mart_order_revenue"}]}, + "assets": [{"type": "table", "name": "mart_order_revenue"}], + } + + assert list(deployer.product_asset_entries(product)) == [ + {"type": "table", "fqn": "aws_glue.lakehouse_dev.sales_order_revenue_gold.mart_order_revenue"} + ] + + def test_deploy_seeds_medallion_buckets_at_storage_service_root( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: (tmp_path / "sales").mkdir() (tmp_path / "sales" / "domain.yaml").write_text( - """name: sales + """apiVersion: openlakeforge.io/v1alpha1 +kind: Domain +name: sales +displayName: Sales +description: Sales domain +status: active data_products: - - name: sales_order_revenue + - id: sales_order_revenue + name: sales_order_revenue + displayName: Sales Order Revenue + description: Revenue from orders. + status: active bronze: - name: raw_orders path: s3://lakehouse-bronze/sales/order_revenue/orders.csv From 1334ac4df055ee55eaf4ce89fb33152b844e9609 Mon Sep 17 00:00:00 2001 From: "a.metwalli" Date: Sun, 19 Jul 2026 17:48:23 +0200 Subject: [PATCH 3/9] fix(contracts): validate domain table groups --- docs/schema/domain.schema.json | 7 ++++++- tools/olf/olf/descriptors.py | 14 +++++++++++++- tools/olf/tests/test_descriptors.py | 27 +++++++++++++++++++++++++++ 3 files changed, 46 insertions(+), 2 deletions(-) diff --git a/docs/schema/domain.schema.json b/docs/schema/domain.schema.json index 80bf098..d349e48 100644 --- a/docs/schema/domain.schema.json +++ b/docs/schema/domain.schema.json @@ -45,7 +45,12 @@ "tableGroup": { "type": "object", "required": ["tables"], - "properties": {"tables": {"type": "array"}}, + "properties": { + "tables": { + "type": "array", + "items": {"type": "object", "required": ["name"]} + } + }, "not": {"required": ["schema"]} } } diff --git a/tools/olf/olf/descriptors.py b/tools/olf/olf/descriptors.py index 58240f0..250abad 100644 --- a/tools/olf/olf/descriptors.py +++ b/tools/olf/olf/descriptors.py @@ -36,7 +36,19 @@ def validate_domain_descriptor(document: Mapping[str, Any], *, source: str = "do if field not in product: raise DomainDescriptorError(f"{source}: data_products[{index}] missing {field!r}") for group in ("silver_tables", "gold_tables"): - spec = product.get(group) or {} + if group not in product: + continue + spec = product[group] + if not isinstance(spec, Mapping): + raise DomainDescriptorError(f"{source}: data_products[{index}].{group} must be an object") + tables = spec.get("tables") + if not isinstance(tables, list): + raise DomainDescriptorError(f"{source}: data_products[{index}].{group}.tables must be an array") + for table_index, table in enumerate(tables): + if not isinstance(table, Mapping) or not table.get("name"): + raise DomainDescriptorError( + f"{source}: data_products[{index}].{group}.tables[{table_index}] must have a name" + ) if "schema" in spec: raise DomainDescriptorError(f"{source}: {group}.schema must be derived from provider contracts") for asset_index, asset in enumerate(product.get("assets") or []): diff --git a/tools/olf/tests/test_descriptors.py b/tools/olf/tests/test_descriptors.py index dbda16b..5c4d6e1 100644 --- a/tools/olf/tests/test_descriptors.py +++ b/tools/olf/tests/test_descriptors.py @@ -63,3 +63,30 @@ def test_domain_descriptor_rejects_legacy_string_asset() -> None: } with pytest.raises(DomainDescriptorError, match="logical asset object"): validate_domain_descriptor(descriptor) + + +@pytest.mark.parametrize( + ("group", "value"), + [("silver_tables", {}), ("gold_tables", {"tables": ["mart_revenue"]})], +) +def test_domain_descriptor_rejects_malformed_table_groups(group: str, value: object) -> None: + descriptor = { + "apiVersion": "openlakeforge.io/v1alpha1", + "kind": "Domain", + "name": "sales", + "displayName": "Sales", + "description": "Sales", + "status": "planned", + "data_products": [ + { + "id": "orders", + "name": "sales_orders", + "displayName": "Orders", + "description": "Orders", + "status": "planned", + group: value, + } + ], + } + with pytest.raises(DomainDescriptorError): + validate_domain_descriptor(descriptor) From 74e30f3aaa29670a0b3b185b34acff0475be8b2f Mon Sep 17 00:00:00 2001 From: "a.metwalli" Date: Sun, 19 Jul 2026 17:57:55 +0200 Subject: [PATCH 4/9] fix: require logical asset names --- docs/schema/domain.schema.json | 4 ++++ tools/olf/olf/descriptors.py | 4 ++++ tools/olf/tests/test_descriptors.py | 23 +++++++++++++++++++++++ 3 files changed, 31 insertions(+) diff --git a/docs/schema/domain.schema.json b/docs/schema/domain.schema.json index d349e48..48d397b 100644 --- a/docs/schema/domain.schema.json +++ b/docs/schema/domain.schema.json @@ -29,6 +29,10 @@ "type": "array", "items": { "type": "object", + "required": ["name"], + "properties": { + "name": {"type": "string", "minLength": 1} + }, "not": { "anyOf": [ {"required": ["fqn"]}, diff --git a/tools/olf/olf/descriptors.py b/tools/olf/olf/descriptors.py index 250abad..bf0f150 100644 --- a/tools/olf/olf/descriptors.py +++ b/tools/olf/olf/descriptors.py @@ -64,6 +64,10 @@ def validate_domain_descriptor(document: Mapping[str, Any], *, source: str = "do raise DomainDescriptorError( f"{source}: data_products[{index}].assets[{asset_index}] must not contain physical FQNs" ) + if not asset.get("name"): + raise DomainDescriptorError( + f"{source}: data_products[{index}].assets[{asset_index}] must have a logical name" + ) def load_domain_descriptor(path: str | Path) -> dict[str, Any]: diff --git a/tools/olf/tests/test_descriptors.py b/tools/olf/tests/test_descriptors.py index 5c4d6e1..0fbce74 100644 --- a/tools/olf/tests/test_descriptors.py +++ b/tools/olf/tests/test_descriptors.py @@ -65,6 +65,29 @@ def test_domain_descriptor_rejects_legacy_string_asset() -> None: validate_domain_descriptor(descriptor) +def test_domain_descriptor_rejects_logical_asset_without_name() -> None: + descriptor = { + "apiVersion": "openlakeforge.io/v1alpha1", + "kind": "Domain", + "name": "sales", + "displayName": "Sales", + "description": "Sales", + "status": "planned", + "data_products": [ + { + "id": "orders", + "name": "sales_orders", + "displayName": "Orders", + "description": "Orders", + "status": "planned", + "assets": [{"type": "table"}], + } + ], + } + with pytest.raises(DomainDescriptorError, match="logical name"): + validate_domain_descriptor(descriptor) + + @pytest.mark.parametrize( ("group", "value"), [("silver_tables", {}), ("gold_tables", {"tables": ["mart_revenue"]})], From f43e70049aaa9f4fbef4a0c339e4e4469bfd8ac4 Mon Sep 17 00:00:00 2001 From: "a.metwalli" Date: Sun, 19 Jul 2026 18:07:06 +0200 Subject: [PATCH 5/9] fix: reject physical identities in table entries --- docs/schema/domain.schema.json | 11 ++++++++++- tools/olf/olf/descriptors.py | 5 +++++ tools/olf/tests/test_descriptors.py | 24 ++++++++++++++++++++++++ 3 files changed, 39 insertions(+), 1 deletion(-) diff --git a/docs/schema/domain.schema.json b/docs/schema/domain.schema.json index 48d397b..5b68ee2 100644 --- a/docs/schema/domain.schema.json +++ b/docs/schema/domain.schema.json @@ -52,7 +52,16 @@ "properties": { "tables": { "type": "array", - "items": {"type": "object", "required": ["name"]} + "items": { + "type": "object", + "required": ["name"], + "not": { + "anyOf": [ + {"required": ["fqn"]}, + {"required": ["fullyQualifiedName"]} + ] + } + } } }, "not": {"required": ["schema"]} diff --git a/tools/olf/olf/descriptors.py b/tools/olf/olf/descriptors.py index bf0f150..02fe608 100644 --- a/tools/olf/olf/descriptors.py +++ b/tools/olf/olf/descriptors.py @@ -49,6 +49,11 @@ def validate_domain_descriptor(document: Mapping[str, Any], *, source: str = "do raise DomainDescriptorError( f"{source}: data_products[{index}].{group}.tables[{table_index}] must have a name" ) + if "fqn" in table or "fullyQualifiedName" in table: + raise DomainDescriptorError( + f"{source}: data_products[{index}].{group}.tables[{table_index}] " + "must not contain physical FQNs" + ) if "schema" in spec: raise DomainDescriptorError(f"{source}: {group}.schema must be derived from provider contracts") for asset_index, asset in enumerate(product.get("assets") or []): diff --git a/tools/olf/tests/test_descriptors.py b/tools/olf/tests/test_descriptors.py index 0fbce74..d8fdd00 100644 --- a/tools/olf/tests/test_descriptors.py +++ b/tools/olf/tests/test_descriptors.py @@ -88,6 +88,30 @@ def test_domain_descriptor_rejects_logical_asset_without_name() -> None: validate_domain_descriptor(descriptor) +@pytest.mark.parametrize("field", ["fqn", "fullyQualifiedName"]) +def test_domain_descriptor_rejects_physical_table_identity(field: str) -> None: + descriptor = { + "apiVersion": "openlakeforge.io/v1alpha1", + "kind": "Domain", + "name": "sales", + "displayName": "Sales", + "description": "Sales", + "status": "planned", + "data_products": [ + { + "id": "orders", + "name": "sales_orders", + "displayName": "Orders", + "description": "Orders", + "status": "planned", + "gold_tables": {"tables": [{"name": "mart_orders", field: "polaris.lakehouse.mart_orders"}]}, + } + ], + } + with pytest.raises(DomainDescriptorError, match="physical FQNs"): + validate_domain_descriptor(descriptor) + + @pytest.mark.parametrize( ("group", "value"), [("silver_tables", {}), ("gold_tables", {"tables": ["mart_revenue"]})], From b5c0a344e7241b42b38a6fceabfa017b268399e8 Mon Sep 17 00:00:00 2001 From: "a.metwalli" Date: Sun, 19 Jul 2026 18:14:46 +0200 Subject: [PATCH 6/9] fix: validate logical table metadata --- docs/schema/domain.schema.json | 6 +++- tools/olf/olf/descriptors.py | 13 +++++--- tools/olf/tests/test_descriptors.py | 46 +++++++++++++++++++++++++++++ 3 files changed, 60 insertions(+), 5 deletions(-) diff --git a/docs/schema/domain.schema.json b/docs/schema/domain.schema.json index 5b68ee2..281a8db 100644 --- a/docs/schema/domain.schema.json +++ b/docs/schema/domain.schema.json @@ -31,7 +31,8 @@ "type": "object", "required": ["name"], "properties": { - "name": {"type": "string", "minLength": 1} + "name": {"type": "string", "minLength": 1}, + "type": {"const": "table"} }, "not": { "anyOf": [ @@ -55,6 +56,9 @@ "items": { "type": "object", "required": ["name"], + "properties": { + "name": {"type": "string", "minLength": 1} + }, "not": { "anyOf": [ {"required": ["fqn"]}, diff --git a/tools/olf/olf/descriptors.py b/tools/olf/olf/descriptors.py index 02fe608..e7e0d63 100644 --- a/tools/olf/olf/descriptors.py +++ b/tools/olf/olf/descriptors.py @@ -45,9 +45,10 @@ def validate_domain_descriptor(document: Mapping[str, Any], *, source: str = "do if not isinstance(tables, list): raise DomainDescriptorError(f"{source}: data_products[{index}].{group}.tables must be an array") for table_index, table in enumerate(tables): - if not isinstance(table, Mapping) or not table.get("name"): + if not isinstance(table, Mapping) or not isinstance(table.get("name"), str) or not table["name"]: raise DomainDescriptorError( - f"{source}: data_products[{index}].{group}.tables[{table_index}] must have a name" + f"{source}: data_products[{index}].{group}.tables[{table_index}] " + "must have a non-empty string name" ) if "fqn" in table or "fullyQualifiedName" in table: raise DomainDescriptorError( @@ -69,9 +70,13 @@ def validate_domain_descriptor(document: Mapping[str, Any], *, source: str = "do raise DomainDescriptorError( f"{source}: data_products[{index}].assets[{asset_index}] must not contain physical FQNs" ) - if not asset.get("name"): + if not isinstance(asset.get("name"), str) or not asset["name"]: + raise DomainDescriptorError( + f"{source}: data_products[{index}].assets[{asset_index}] must have a non-empty logical name" + ) + if asset.get("type") not in (None, "table"): raise DomainDescriptorError( - f"{source}: data_products[{index}].assets[{asset_index}] must have a logical name" + f"{source}: data_products[{index}].assets[{asset_index}] must have type 'table' when specified" ) diff --git a/tools/olf/tests/test_descriptors.py b/tools/olf/tests/test_descriptors.py index d8fdd00..d83186b 100644 --- a/tools/olf/tests/test_descriptors.py +++ b/tools/olf/tests/test_descriptors.py @@ -112,6 +112,52 @@ def test_domain_descriptor_rejects_physical_table_identity(field: str) -> None: validate_domain_descriptor(descriptor) +def test_domain_descriptor_rejects_non_string_table_name() -> None: + descriptor = { + "apiVersion": "openlakeforge.io/v1alpha1", + "kind": "Domain", + "name": "sales", + "displayName": "Sales", + "description": "Sales", + "status": "planned", + "data_products": [ + { + "id": "orders", + "name": "sales_orders", + "displayName": "Orders", + "description": "Orders", + "status": "planned", + "gold_tables": {"tables": [{"name": 123}]}, + } + ], + } + with pytest.raises(DomainDescriptorError, match="non-empty string name"): + validate_domain_descriptor(descriptor) + + +def test_domain_descriptor_rejects_unsupported_logical_asset_type() -> None: + descriptor = { + "apiVersion": "openlakeforge.io/v1alpha1", + "kind": "Domain", + "name": "sales", + "displayName": "Sales", + "description": "Sales", + "status": "planned", + "data_products": [ + { + "id": "orders", + "name": "sales_orders", + "displayName": "Orders", + "description": "Orders", + "status": "planned", + "assets": [{"type": "dashboard", "name": "mart_orders"}], + } + ], + } + with pytest.raises(DomainDescriptorError, match="type 'table'"): + validate_domain_descriptor(descriptor) + + @pytest.mark.parametrize( ("group", "value"), [("silver_tables", {}), ("gold_tables", {"tables": ["mart_revenue"]})], From 265325650f1572ac6f8a125038aebd7f73bfd6ec Mon Sep 17 00:00:00 2001 From: "a.metwalli" Date: Sun, 19 Jul 2026 18:22:39 +0200 Subject: [PATCH 7/9] fix: preflight provider schema coverage --- tools/olf/olf/openmetadata.py | 20 +++++++++++----- tools/olf/tests/test_openmetadata.py | 34 ++++++++++++++++++++++++++++ 2 files changed, 48 insertions(+), 6 deletions(-) diff --git a/tools/olf/olf/openmetadata.py b/tools/olf/olf/openmetadata.py index aebbc43..578d96d 100644 --- a/tools/olf/olf/openmetadata.py +++ b/tools/olf/olf/openmetadata.py @@ -235,27 +235,34 @@ def __init__(self, config: OpenMetadataConfig, client: OpenMetadataClient): # --- product/table spec helpers --------------------------------------- - def schema_fqn_for_product(self, product: dict, table_group_key: str, fallback): + def schema_fqn_for_product(self, product: dict, table_group_key: str) -> str | None: product_key = product_contract_key(product) if table_group_key == "silver_tables": - return self.config.catalog_silver_schema_fqns.get(product_key) or fallback + return self.config.catalog_silver_schema_fqns.get(product_key) if table_group_key == "gold_tables": - return self.config.catalog_gold_schema_fqns.get(product_key) or fallback - return fallback + return self.config.catalog_gold_schema_fqns.get(product_key) + return None def product_table_specs(self, product: dict): for key in ["silver_tables", "gold_tables"]: spec = product.get(key) if not spec: continue - schema_fqn = self.schema_fqn_for_product(product, key, spec.get("schema")) + schema_fqn = self.schema_fqn_for_product(product, key) if not schema_fqn: raise OpenMetadataError( - f"Data product '{product.get('name')}' table group '{key}' is missing required 'schema'." + f"Data product '{product_contract_key(product)}' table group '{key}' is not covered " + "by the provider contract schema FQNs." ) for table in spec.get("tables", []): yield schema_fqn, table + def validate_provider_schema_coverage(self, domain_specs: list[tuple[Path, dict]]) -> None: + """Fail before metadata writes when descriptors outpace provider contract namespaces.""" + for _, domain in domain_specs: + for product in product_entries(domain): + list(self.product_table_specs(product)) + def provider_asset_fqn(self, product: dict, fqn): if not fqn: return fqn @@ -571,6 +578,7 @@ def deploy(self) -> None: raise OpenMetadataError( f"No OpenMetadata domain metadata files found under {self.config.metadata_root}//domain.yaml" ) + self.validate_provider_schema_coverage(domain_specs) # Phase A+B: Object Store service and medallion bucket containers. self.ensure_storage_service() diff --git a/tools/olf/tests/test_openmetadata.py b/tools/olf/tests/test_openmetadata.py index f898e7f..aa3264c 100644 --- a/tools/olf/tests/test_openmetadata.py +++ b/tools/olf/tests/test_openmetadata.py @@ -177,6 +177,40 @@ def test_logical_asset_name_resolves_through_provider_contract() -> None: ] +def test_provider_schema_coverage_rejects_new_product_without_contract() -> None: + cfg = om.OpenMetadataConfig.from_environment( + {}, + base_url="http://x", + admin_email="a", + admin_password="p", + metadata_root="domains", + metadata_source_dir="", + allow_missing_assets=False, + catalog_service="polaris", + catalog_database="lakehouse_dev", + cleanup_legacy_default_database=False, + ) + deployer = om.OpenMetadataDeployer(cfg, om.OpenMetadataClient(cfg.base_url)) + domain_specs = [ + ( + Path("domains/sales/domain.yaml"), + { + "name": "sales", + "data_products": [ + { + "id": "new_product", + "name": "sales_new_product", + "gold_tables": {"tables": [{"name": "mart_new_product"}]}, + } + ], + }, + ) + ] + + with pytest.raises(om.OpenMetadataError, match="not covered by the provider contract"): + deployer.validate_provider_schema_coverage(domain_specs) + + def test_deploy_seeds_medallion_buckets_at_storage_service_root( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: From 3e94ed4d592baafa25c915a296a9bd2044b045ff Mon Sep 17 00:00:00 2001 From: "a.metwalli" Date: Sun, 19 Jul 2026 18:30:28 +0200 Subject: [PATCH 8/9] fix: default direct metadata schema contracts --- tools/olf/olf/openmetadata.py | 29 ++++++++++++++---- tools/olf/tests/test_openmetadata.py | 45 ++++++++++++++++++++++++++++ 2 files changed, 68 insertions(+), 6 deletions(-) diff --git a/tools/olf/olf/openmetadata.py b/tools/olf/olf/openmetadata.py index 578d96d..57ee2ef 100644 --- a/tools/olf/olf/openmetadata.py +++ b/tools/olf/olf/openmetadata.py @@ -20,6 +20,12 @@ from olf.descriptors import load_domain_descriptor +_SEED_PRODUCT_KEYS = ( + "sales_order_revenue", + "sales_customer_health", + "supply_chain_inventory_reliability", +) + class OpenMetadataError(RuntimeError): pass @@ -62,9 +68,13 @@ def from_environment( catalog_database: str, cleanup_legacy_default_database: bool, ) -> OpenMetadataConfig: + catalog_service = catalog_service or "polaris" + catalog_database = catalog_database or "lakehouse_dev" catalog_database_fqn = environ.get( "OPENLAKEFORGE_CATALOG_DATABASE_FQN", f"{catalog_service}.{catalog_database}" ) + silver_schema_fqns_raw = environ.get("OPENLAKEFORGE_CATALOG_SILVER_SCHEMA_FQNS_JSON") + gold_schema_fqns_raw = environ.get("OPENLAKEFORGE_CATALOG_GOLD_SCHEMA_FQNS_JSON") return cls( base_url=base_url.rstrip("/"), admin_email=admin_email, @@ -76,13 +86,15 @@ def from_environment( catalog_database=catalog_database, cleanup_legacy_default_database=cleanup_legacy_default_database, catalog_database_fqn=catalog_database_fqn, - catalog_silver_schema_fqns=_parse_json_env( - "OPENLAKEFORGE_CATALOG_SILVER_SCHEMA_FQNS_JSON", - environ.get("OPENLAKEFORGE_CATALOG_SILVER_SCHEMA_FQNS_JSON", "{}"), + catalog_silver_schema_fqns=( + _parse_json_env("OPENLAKEFORGE_CATALOG_SILVER_SCHEMA_FQNS_JSON", silver_schema_fqns_raw) + if silver_schema_fqns_raw + else _default_schema_fqns(catalog_database_fqn, "silver") ), - catalog_gold_schema_fqns=_parse_json_env( - "OPENLAKEFORGE_CATALOG_GOLD_SCHEMA_FQNS_JSON", - environ.get("OPENLAKEFORGE_CATALOG_GOLD_SCHEMA_FQNS_JSON", "{}"), + catalog_gold_schema_fqns=( + _parse_json_env("OPENLAKEFORGE_CATALOG_GOLD_SCHEMA_FQNS_JSON", gold_schema_fqns_raw) + if gold_schema_fqns_raw + else _default_schema_fqns(catalog_database_fqn, "gold") ), storage_service=environ.get("OPENLAKEFORGE_STORAGE_OM_SERVICE", "seaweedfs"), storage_display_name=environ.get("OPENLAKEFORGE_STORAGE_DISPLAY_NAME", "SeaweedFS S3"), @@ -107,6 +119,11 @@ def _parse_json_env(name: str, raw: str) -> dict: return value +def _default_schema_fqns(catalog_database_fqn: str, layer: str) -> dict[str, str]: + """Return the seed-product contract used by direct local CLI execution.""" + return {product: f"{catalog_database_fqn}.{product}_{layer}" for product in _SEED_PRODUCT_KEYS} + + @dataclass class OpenMetadataClient: base_url: str diff --git a/tools/olf/tests/test_openmetadata.py b/tools/olf/tests/test_openmetadata.py index aa3264c..975745b 100644 --- a/tools/olf/tests/test_openmetadata.py +++ b/tools/olf/tests/test_openmetadata.py @@ -77,6 +77,51 @@ def test_config_from_environment_reads_schema_fqns() -> None: assert cfg.storage_gold_bucket == "openlakeforge-poc-gold" +def test_config_from_environment_defaults_seed_schema_fqns_for_direct_cli() -> None: + cfg = om.OpenMetadataConfig.from_environment( + {}, + base_url="http://x", + admin_email="a", + admin_password="p", + metadata_root="domains", + metadata_source_dir="", + allow_missing_assets=False, + catalog_service="", + catalog_database="", + cleanup_legacy_default_database=False, + ) + + assert cfg.catalog_service == "polaris" + assert cfg.catalog_database == "lakehouse_dev" + assert cfg.catalog_silver_schema_fqns["sales_order_revenue"] == ( + "polaris.lakehouse_dev.sales_order_revenue_silver" + ) + assert cfg.catalog_gold_schema_fqns["sales_order_revenue"] == ( + "polaris.lakehouse_dev.sales_order_revenue_gold" + ) + + +def test_config_from_environment_preserves_explicit_empty_schema_contract() -> None: + cfg = om.OpenMetadataConfig.from_environment( + { + "OPENLAKEFORGE_CATALOG_SILVER_SCHEMA_FQNS_JSON": "{}", + "OPENLAKEFORGE_CATALOG_GOLD_SCHEMA_FQNS_JSON": "{}", + }, + base_url="http://x", + admin_email="a", + admin_password="p", + metadata_root="domains", + metadata_source_dir="", + allow_missing_assets=False, + catalog_service="polaris", + catalog_database="lakehouse_dev", + cleanup_legacy_default_database=False, + ) + + assert cfg.catalog_silver_schema_fqns == {} + assert cfg.catalog_gold_schema_fqns == {} + + def test_storage_bucket_specs_dedup() -> None: cfg = om.OpenMetadataConfig.from_environment( {}, From e5123ab8d40459d043e2959f2b0c2b37b2ba5f2a Mon Sep 17 00:00:00 2001 From: "a.metwalli" Date: Sun, 19 Jul 2026 18:39:42 +0200 Subject: [PATCH 9/9] fix: preflight descriptor deployment inputs --- docs/schema/domain.schema.json | 33 ++++++++++--- tools/olf/olf/descriptors.py | 55 +++++++++++++++++++++ tools/olf/olf/openmetadata.py | 33 +++++++++---- tools/olf/tests/test_descriptors.py | 23 +++++++++ tools/olf/tests/test_openmetadata.py | 71 +++++++++++++++++++++++++++- 5 files changed, 199 insertions(+), 16 deletions(-) diff --git a/docs/schema/domain.schema.json b/docs/schema/domain.schema.json index 281a8db..24540c9 100644 --- a/docs/schema/domain.schema.json +++ b/docs/schema/domain.schema.json @@ -11,18 +11,29 @@ "name": {"type": "string", "pattern": "^[a-z][a-z0-9_]*$"}, "displayName": {"type": "string", "minLength": 1}, "description": {"type": "string"}, - "status": {"type": "string"}, + "status": {"type": "string", "minLength": 1}, "data_products": { "type": "array", "items": { "type": "object", "required": ["id", "name", "displayName", "description", "status"], "properties": { - "id": {"type": "string"}, - "name": {"type": "string"}, - "displayName": {"type": "string"}, + "id": {"type": "string", "minLength": 1}, + "name": {"type": "string", "minLength": 1}, + "displayName": {"type": "string", "minLength": 1}, "description": {"type": "string"}, - "status": {"type": "string"}, + "status": {"type": "string", "minLength": 1}, + "asset_prefix": {"type": "string", "minLength": 1}, + "domain": {"type": "string", "minLength": 1}, + "domains": { + "type": "array", + "minItems": 1, + "items": {"type": "string", "minLength": 1} + }, + "bronze": { + "type": "array", + "items": {"$ref": "#/$defs/bronzeEntry"} + }, "silver_tables": {"$ref": "#/$defs/tableGroup"}, "gold_tables": {"$ref": "#/$defs/tableGroup"}, "assets": { @@ -57,7 +68,8 @@ "type": "object", "required": ["name"], "properties": { - "name": {"type": "string", "minLength": 1} + "name": {"type": "string", "minLength": 1}, + "description": {"type": "string"} }, "not": { "anyOf": [ @@ -69,6 +81,15 @@ } }, "not": {"required": ["schema"]} + }, + "bronzeEntry": { + "type": "object", + "required": ["name", "path"], + "properties": { + "name": {"type": "string", "minLength": 1}, + "path": {"type": "string", "minLength": 1}, + "description": {"type": "string"} + } } } } diff --git a/tools/olf/olf/descriptors.py b/tools/olf/olf/descriptors.py index e7e0d63..4768354 100644 --- a/tools/olf/olf/descriptors.py +++ b/tools/olf/olf/descriptors.py @@ -2,6 +2,7 @@ from __future__ import annotations +import re from collections.abc import Mapping from pathlib import Path from typing import Any @@ -27,6 +28,13 @@ def validate_domain_descriptor(document: Mapping[str, Any], *, source: str = "do for field in ("name", "displayName", "description", "status", "data_products"): if field not in document: raise DomainDescriptorError(f"{source}: missing required field {field!r}") + if not isinstance(document["name"], str) or not re.fullmatch(r"[a-z][a-z0-9_]*", document["name"]): + raise DomainDescriptorError(f"{source}: name must match '^[a-z][a-z0-9_]*$'") + for field in ("displayName", "status"): + if not isinstance(document[field], str) or not document[field]: + raise DomainDescriptorError(f"{source}: {field} must be a non-empty string") + if not isinstance(document["description"], str): + raise DomainDescriptorError(f"{source}: description must be a string") if not isinstance(document["data_products"], list): raise DomainDescriptorError(f"{source}: data_products must be an array") for index, product in enumerate(document["data_products"]): @@ -35,6 +43,46 @@ def validate_domain_descriptor(document: Mapping[str, Any], *, source: str = "do for field in ("id", "name", "displayName", "description", "status"): if field not in product: raise DomainDescriptorError(f"{source}: data_products[{index}] missing {field!r}") + for field in ("id", "name", "displayName", "status"): + if not isinstance(product[field], str) or not product[field]: + raise DomainDescriptorError( + f"{source}: data_products[{index}].{field} must be a non-empty string" + ) + if not isinstance(product["description"], str): + raise DomainDescriptorError(f"{source}: data_products[{index}].description must be a string") + if "asset_prefix" in product and ( + not isinstance(product["asset_prefix"], str) or not product["asset_prefix"] + ): + raise DomainDescriptorError(f"{source}: data_products[{index}].asset_prefix must be a non-empty string") + if "domain" in product and (not isinstance(product["domain"], str) or not product["domain"]): + raise DomainDescriptorError(f"{source}: data_products[{index}].domain must be a non-empty string") + if "domains" in product and ( + not isinstance(product["domains"], list) + or not product["domains"] + or any(not isinstance(domain, str) or not domain for domain in product["domains"]) + ): + raise DomainDescriptorError( + f"{source}: data_products[{index}].domains must be a non-empty array of strings" + ) + if "bronze" in product: + bronze_entries = product["bronze"] + if not isinstance(bronze_entries, list): + raise DomainDescriptorError(f"{source}: data_products[{index}].bronze must be an array") + for bronze_index, bronze_entry in enumerate(bronze_entries): + if not isinstance(bronze_entry, Mapping): + raise DomainDescriptorError( + f"{source}: data_products[{index}].bronze[{bronze_index}] must be an object" + ) + for field in ("name", "path"): + if not isinstance(bronze_entry.get(field), str) or not bronze_entry[field]: + raise DomainDescriptorError( + f"{source}: data_products[{index}].bronze[{bronze_index}].{field} " + "must be a non-empty string" + ) + if "description" in bronze_entry and not isinstance(bronze_entry["description"], str): + raise DomainDescriptorError( + f"{source}: data_products[{index}].bronze[{bronze_index}].description must be a string" + ) for group in ("silver_tables", "gold_tables"): if group not in product: continue @@ -55,8 +103,15 @@ def validate_domain_descriptor(document: Mapping[str, Any], *, source: str = "do f"{source}: data_products[{index}].{group}.tables[{table_index}] " "must not contain physical FQNs" ) + if "description" in table and not isinstance(table["description"], str): + raise DomainDescriptorError( + f"{source}: data_products[{index}].{group}.tables[{table_index}].description " + "must be a string" + ) if "schema" in spec: raise DomainDescriptorError(f"{source}: {group}.schema must be derived from provider contracts") + if "assets" in product and not isinstance(product["assets"], list): + raise DomainDescriptorError(f"{source}: data_products[{index}].assets must be an array") for asset_index, asset in enumerate(product.get("assets") or []): if isinstance(asset, str): raise DomainDescriptorError( diff --git a/tools/olf/olf/openmetadata.py b/tools/olf/olf/openmetadata.py index 57ee2ef..4fdaaae 100644 --- a/tools/olf/olf/openmetadata.py +++ b/tools/olf/olf/openmetadata.py @@ -274,11 +274,12 @@ def product_table_specs(self, product: dict): for table in spec.get("tables", []): yield schema_fqn, table - def validate_provider_schema_coverage(self, domain_specs: list[tuple[Path, dict]]) -> None: - """Fail before metadata writes when descriptors outpace provider contract namespaces.""" + def validate_deployment_inputs(self, domain_specs: list[tuple[Path, dict]]) -> None: + """Resolve every declared table and logical asset before metadata writes.""" for _, domain in domain_specs: for product in product_entries(domain): - list(self.product_table_specs(product)) + list(self.product_asset_entries(product)) + self.validate_bronze_entries(domain_specs) def provider_asset_fqn(self, product: dict, fqn): if not fqn: @@ -541,10 +542,27 @@ def ensure_table_stub(self, schema_fqn, name, description) -> None: def validate_bronze_entries(self, domain_specs) -> None: for _, domain in domain_specs: for product in product_entries(domain): - for container in product.get("bronze") or []: - if not container.get("path"): + bronze_entries = product.get("bronze") + if bronze_entries is None: + continue + if not isinstance(bronze_entries, list): + raise OpenMetadataError( + f"Data product '{product['name']}' Bronze entries must be an array." + ) + for index, container in enumerate(bronze_entries): + if not isinstance(container, dict): + raise OpenMetadataError( + f"Data product '{product['name']}' Bronze entry at index {index} must be an object." + ) + if not isinstance(container.get("name"), str) or not container["name"]: raise OpenMetadataError( - f"Data product '{product['name']}' Bronze entry is missing required 'path'." + f"Data product '{product['name']}' Bronze entry at index {index} " + "is missing required 'name'." + ) + if not isinstance(container.get("path"), str) or not container["path"]: + raise OpenMetadataError( + f"Data product '{product['name']}' Bronze entry at index {index} " + "is missing required 'path'." ) def cleanup_legacy_default_database(self) -> None: @@ -595,11 +613,10 @@ def deploy(self) -> None: raise OpenMetadataError( f"No OpenMetadata domain metadata files found under {self.config.metadata_root}//domain.yaml" ) - self.validate_provider_schema_coverage(domain_specs) + self.validate_deployment_inputs(domain_specs) # Phase A+B: Object Store service and medallion bucket containers. self.ensure_storage_service() - self.validate_bronze_entries(domain_specs) for container in self.storage_bucket_specs(): self.ensure_container(container["name"], None, container["path"], container["description"]) diff --git a/tools/olf/tests/test_descriptors.py b/tools/olf/tests/test_descriptors.py index d83186b..983c1d0 100644 --- a/tools/olf/tests/test_descriptors.py +++ b/tools/olf/tests/test_descriptors.py @@ -135,6 +135,29 @@ def test_domain_descriptor_rejects_non_string_table_name() -> None: validate_domain_descriptor(descriptor) +def test_domain_descriptor_rejects_malformed_bronze_entry() -> None: + descriptor = { + "apiVersion": "openlakeforge.io/v1alpha1", + "kind": "Domain", + "name": "sales", + "displayName": "Sales", + "description": "Sales", + "status": "planned", + "data_products": [ + { + "id": "orders", + "name": "sales_orders", + "displayName": "Orders", + "description": "Orders", + "status": "planned", + "bronze": ["raw_orders"], + } + ], + } + with pytest.raises(DomainDescriptorError, match="bronze\\[0\\] must be an object"): + validate_domain_descriptor(descriptor) + + def test_domain_descriptor_rejects_unsupported_logical_asset_type() -> None: descriptor = { "apiVersion": "openlakeforge.io/v1alpha1", diff --git a/tools/olf/tests/test_openmetadata.py b/tools/olf/tests/test_openmetadata.py index 975745b..a1a1c87 100644 --- a/tools/olf/tests/test_openmetadata.py +++ b/tools/olf/tests/test_openmetadata.py @@ -222,7 +222,7 @@ def test_logical_asset_name_resolves_through_provider_contract() -> None: ] -def test_provider_schema_coverage_rejects_new_product_without_contract() -> None: +def test_deployment_input_validation_rejects_new_product_without_contract() -> None: cfg = om.OpenMetadataConfig.from_environment( {}, base_url="http://x", @@ -253,7 +253,74 @@ def test_provider_schema_coverage_rejects_new_product_without_contract() -> None ] with pytest.raises(om.OpenMetadataError, match="not covered by the provider contract"): - deployer.validate_provider_schema_coverage(domain_specs) + deployer.validate_deployment_inputs(domain_specs) + + +def test_deployment_input_validation_rejects_unknown_logical_asset() -> None: + cfg = om.OpenMetadataConfig.from_environment( + { + "OPENLAKEFORGE_CATALOG_GOLD_SCHEMA_FQNS_JSON": ( + '{"sales_order_revenue": "polaris.lakehouse_dev.sales_order_revenue_gold"}' + ) + }, + base_url="http://x", + admin_email="a", + admin_password="p", + metadata_root="domains", + metadata_source_dir="", + allow_missing_assets=False, + catalog_service="polaris", + catalog_database="lakehouse_dev", + cleanup_legacy_default_database=False, + ) + deployer = om.OpenMetadataDeployer(cfg, om.OpenMetadataClient(cfg.base_url)) + domain_specs = [ + ( + Path("domains/sales/domain.yaml"), + { + "name": "sales", + "data_products": [ + { + "id": "sales_order_revenue", + "name": "sales_order_revenue", + "gold_tables": {"tables": [{"name": "mart_order_revenue"}]}, + "assets": [{"type": "table", "name": "mart_typo"}], + } + ], + }, + ) + ] + + with pytest.raises(om.OpenMetadataError, match="not declared in the product table contract"): + deployer.validate_deployment_inputs(domain_specs) + + +def test_deployment_input_validation_rejects_malformed_bronze_before_writes() -> None: + cfg = om.OpenMetadataConfig.from_environment( + {}, + base_url="http://x", + admin_email="a", + admin_password="p", + metadata_root="domains", + metadata_source_dir="", + allow_missing_assets=False, + catalog_service="polaris", + catalog_database="lakehouse_dev", + cleanup_legacy_default_database=False, + ) + deployer = om.OpenMetadataDeployer(cfg, om.OpenMetadataClient(cfg.base_url)) + domain_specs = [ + ( + Path("domains/sales/domain.yaml"), + { + "name": "sales", + "data_products": [{"name": "sales_order_revenue", "bronze": ["raw_orders"]}], + }, + ) + ] + + with pytest.raises(om.OpenMetadataError, match="Bronze entry at index 0 must be an object"): + deployer.validate_deployment_inputs(domain_specs) def test_deploy_seeds_medallion_buckets_at_storage_service_root(