91 lines
3.0 KiB
Python
91 lines
3.0 KiB
Python
from __future__ import annotations
|
|
|
|
from collections.abc import Mapping
|
|
from typing import Any
|
|
|
|
|
|
SOURCE_PROVENANCE_METADATA_KEY = "source_provenance"
|
|
SOURCE_REVISION_METADATA_KEY = "source_revision"
|
|
|
|
_PROVENANCE_STRING_FIELDS = {
|
|
"source_type",
|
|
"connector_id",
|
|
"provider",
|
|
"external_id",
|
|
"external_path",
|
|
"external_url",
|
|
"revision",
|
|
"revision_label",
|
|
"observed_at",
|
|
"imported_at",
|
|
}
|
|
|
|
|
|
def normalize_source_revision(value: object) -> str | None:
|
|
if value is None:
|
|
return None
|
|
text = str(value).strip()
|
|
return text or None
|
|
|
|
|
|
def normalize_source_provenance(value: object, *, source_revision: object = None) -> dict[str, Any] | None:
|
|
if value is None:
|
|
revision = normalize_source_revision(source_revision)
|
|
return {"revision": revision} if revision else None
|
|
if not isinstance(value, Mapping):
|
|
raise ValueError("source_provenance must be a JSON object")
|
|
|
|
payload: dict[str, Any] = {}
|
|
for field in _PROVENANCE_STRING_FIELDS:
|
|
normalized = normalize_source_revision(value.get(field))
|
|
if normalized:
|
|
payload[field] = normalized
|
|
|
|
extra = value.get("metadata")
|
|
if isinstance(extra, Mapping) and extra:
|
|
payload["metadata"] = dict(extra)
|
|
elif extra is not None and extra != "":
|
|
raise ValueError("source_provenance.metadata must be a JSON object")
|
|
|
|
revision = normalize_source_revision(source_revision)
|
|
if revision:
|
|
payload["revision"] = revision
|
|
return payload or None
|
|
|
|
|
|
def source_metadata(
|
|
*,
|
|
existing: Mapping[str, Any] | None = None,
|
|
source_provenance: object = None,
|
|
source_revision: object = None,
|
|
) -> dict[str, Any] | None:
|
|
payload = dict(existing or {})
|
|
normalized_provenance = normalize_source_provenance(source_provenance, source_revision=source_revision)
|
|
normalized_revision = normalize_source_revision(source_revision)
|
|
if normalized_provenance:
|
|
payload[SOURCE_PROVENANCE_METADATA_KEY] = normalized_provenance
|
|
if normalized_revision:
|
|
payload[SOURCE_REVISION_METADATA_KEY] = normalized_revision
|
|
elif normalized_provenance and normalized_provenance.get("revision"):
|
|
payload[SOURCE_REVISION_METADATA_KEY] = str(normalized_provenance["revision"])
|
|
return payload or None
|
|
|
|
|
|
def source_provenance_from_metadata(metadata: Mapping[str, Any] | None) -> dict[str, Any] | None:
|
|
if not isinstance(metadata, Mapping):
|
|
return None
|
|
value = metadata.get(SOURCE_PROVENANCE_METADATA_KEY)
|
|
if isinstance(value, Mapping):
|
|
return dict(value)
|
|
return None
|
|
|
|
|
|
def source_revision_from_metadata(metadata: Mapping[str, Any] | None) -> str | None:
|
|
if not isinstance(metadata, Mapping):
|
|
return None
|
|
revision = normalize_source_revision(metadata.get(SOURCE_REVISION_METADATA_KEY))
|
|
if revision:
|
|
return revision
|
|
provenance = source_provenance_from_metadata(metadata)
|
|
return normalize_source_revision(provenance.get("revision") if provenance else None)
|