From 8521fedb396576b22577f9efa1dccc888e1dd4dd Mon Sep 17 00:00:00 2001 From: Mohammed Minhajuddin <68331751+minhajuddin2510@users.noreply.github.com> Date: Mon, 21 Sep 2026 09:45:39 -0400 Subject: [PATCH] feat: mirror CKAN source catalogs exactly --- README.md | 32 ++ docker-compose.yml | 23 +- docker/fairstore/Dockerfile | 11 + scripts/populate_fairstore.py | 741 +++++++++++++++++++------- tests/unit/test_fairstore_populate.py | 202 +++++++ 5 files changed, 813 insertions(+), 196 deletions(-) create mode 100644 docker/fairstore/Dockerfile diff --git a/README.md b/README.md index 01efdbf..af04a05 100644 --- a/README.md +++ b/README.md @@ -358,6 +358,38 @@ examples/ # Sample generated notebooks --- +## CKAN Fair Store mirror + +Start the local CKAN 2.11 Fair Store with organization hierarchy support: + +```bash +docker compose --profile fairstore up -d --build fairstore +``` + +Preview a registered CKAN portal before writing anything: + +```bash +python -m scripts.populate_fairstore \ + --site wprdc \ + --target-url http://localhost:5001 +``` + +To apply the mirror, create a target CKAN sysadmin token, expose it through +`CKAN_API_KEY` (or a protected file), and add `--apply`. The command performs a +collision preflight first, then preserves source organization, dataset, and +resource names and UUIDs. It creates one parent organization for the source +portal, attaches source organizations beneath it, and adds any matching qsv +profile to the resource as separate metadata. Re-running patches the preserved +UUIDs, so it does not duplicate resources. + +```bash +python -m scripts.populate_fairstore \ + --site wprdc \ + --target-url http://localhost:5001 \ + --api-key-file /path/to/protected-token \ + --apply +``` + ## Development ```bash diff --git a/docker-compose.yml b/docker-compose.yml index 881653e..8839665 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -82,7 +82,10 @@ services: # --------------------------------------------------------------------- fairstore: profiles: ["fairstore"] - image: ckan/ckan-dev:2.10 + build: + context: ./docker/fairstore + dockerfile: Dockerfile + image: data-concierge-fairstore:2.11.3 platform: linux/amd64 depends_on: fairstore-db: @@ -90,7 +93,7 @@ services: fairstore-solr: condition: service_started ports: - - "5001:5000" + - "${FAIRSTORE_PORT:-5001}:5000" environment: - CKAN_SITE_URL=http://localhost:5000 # Credentials must match the users the postgres image's init scripts @@ -100,18 +103,12 @@ services: - CKAN_DATASTORE_WRITE_URL=postgresql://ckan_datastore_write:datastore@fairstore-db/datastore - CKAN_DATASTORE_READ_URL=postgresql://ckan_datastore_read:datastore@fairstore-db/datastore - CKAN_SOLR_URL=http://fairstore-solr:8983/solr/ckan - - CKAN__PLUGINS=envvars datastore + # hierarchy_display must load before hierarchy_form. Organizations that + # represent source portals can then be parents of the organizations + # mirrored from those portals. + - CKAN__PLUGINS=envvars datastore hierarchy_display hierarchy_form volumes: - fairstore_data:/var/lib/ckan - # The published ckan-dev:2.10 start script hardcodes the server binary as - # /usr/local/bin/ckan, but the image installs it at /usr/bin/ckan, so the - # dev server exits 127 in a loop. The prerun (config, db init) uses the - # binary via PATH and works; only that final line is wrong. Symlink the - # expected path, then hand off to the stock start script unchanged. - entrypoint: - - /bin/sh - - -c - - "ln -sf /usr/bin/ckan /usr/local/bin/ckan && exec /srv/app/start_ckan_development.sh" restart: unless-stopped fairstore-db: @@ -135,7 +132,7 @@ services: fairstore-solr: profiles: ["fairstore"] - image: ckan/ckan-solr:2.10-solr9 + image: ckan/ckan-solr:2.11-solr9 platform: linux/amd64 volumes: - fairstore_solr:/var/solr diff --git a/docker/fairstore/Dockerfile b/docker/fairstore/Dockerfile new file mode 100644 index 0000000..59a688b --- /dev/null +++ b/docker/fairstore/Dockerfile @@ -0,0 +1,11 @@ +FROM ckan/ckan-dev:2.11.3 + +# Keep the extension revision deterministic and aligned with the CKAN 2.11 +# version used by the hosted Fair Store. This revision's CI covers CKAN 2.11. +ARG CKANEXT_HIERARCHY_SHA=53c1ee74805d79aef7c9c01cc513eb881cb8928f + +USER root +RUN pip install --no-cache-dir \ + "git+https://github.com/ckan/ckanext-hierarchy.git@${CKANEXT_HIERARCHY_SHA}#egg=ckanext-hierarchy" + +USER ckan diff --git a/scripts/populate_fairstore.py b/scripts/populate_fairstore.py index b564335..3338b35 100644 --- a/scripts/populate_fairstore.py +++ b/scripts/populate_fairstore.py @@ -1,38 +1,15 @@ #!/usr/bin/env python -"""Publish onboarded data dictionaries into the Fair Store (CKAN). - -Issue #133. ``scripts/onboard_ckan.py`` builds a rich per-column dictionary -(qsv stats, labels, types, top values) and stores it as ``index.json``. This -takes that dictionary and writes it into the containerized CKAN instance as -real dataset and resource records, so the structured metadata is queryable -through CKAN's deterministic API rather than only through the in-process -search index. - -What lands in CKAN, per resource: - -* a **dataset** (CKAN package) carrying the dataset title, description, - organization and tags; -* a **resource** under it carrying the qsv description and format; -* the **column-level data dictionary** — every column's label, type, stats - and top values — as a JSON blob in the resource's ``extras`` under - ``data_dictionary``, plus flat ``column_count`` / ``row_count`` extras for - filtering. - -The vector store keeps the fuzzy half and points back at these CKAN ids; this -script never touches it. Re-running is idempotent: a dataset is matched by its -slug and updated in place rather than duplicated. - -Usage:: - - # Boot the stack and create a token first: - # docker compose --profile fairstore up -d - # docker compose exec fairstore ckan -c /srv/app/ckan.ini sysadmin add admin - # then generate an API token in the CKAN UI and: - python -m scripts.populate_fairstore \\ - --ckan-url http://localhost:5001 \\ - --api-key \\ - --site wprdc \\ - --org city-of-pittsburgh +"""Mirror a source CKAN catalog into the Fair Store. + +The mirror keeps source organization, group, dataset, and resource names and +UUIDs. A top-level organization represents the source portal; source root +organizations are attached beneath it through ckanext-hierarchy. Available +qsv profiling metadata from ``data/ckan_onboard//index.json`` is added to +the corresponding resource without replacing the source description. + +The command is read-only unless ``--apply`` is supplied. Before any write it +checks the whole source snapshot for name and UUID collisions. Re-running is +idempotent: existing objects are patched by their preserved source UUIDs. """ from __future__ import annotations @@ -40,190 +17,588 @@ import argparse import asyncio import json +import os import re import sys +import uuid +from collections.abc import Iterable from pathlib import Path from typing import Any _PROJECT_ROOT = Path(__file__).resolve().parent.parent sys.path.insert(0, str(_PROJECT_ROOT / "src")) -from data_concierge.core.logging import get_logger # noqa: E402 from data_concierge.data_layer.connectors.ckan import CKANClient # noqa: E402 from data_concierge.data_layer.onboard_index import _scrub_secrets # noqa: E402 - -logger = get_logger("populate_fairstore") +from data_concierge.gateway.ckan_sites import get_site # noqa: E402 _SLUG_RE = re.compile(r"[^a-z0-9-]+") +_PROVENANCE_KEYS = { + "portal": "mirror_source_portal", + "url": "mirror_source_url", + "id": "mirror_source_id", +} + +_ORGANIZATION_FIELDS = ( + "id", + "name", + "title", + "description", + "image_url", + "type", + "state", + "approval_status", +) +_GROUP_FIELDS = ("id", "name", "title", "description", "image_url", "type", "state") +_PACKAGE_FIELDS = ( + "id", + "name", + "title", + "author", + "author_email", + "maintainer", + "maintainer_email", + "license_id", + "notes", + "url", + "version", + "state", + "type", + "private", + "plugin_data", +) +_RESOURCE_FIELDS = ( + "id", + "url", + "description", + "format", + "hash", + "name", + "resource_type", + "url_type", + "mimetype", + "mimetype_inner", + "cache_url", + "size", + "created", + "last_modified", + "cache_last_updated", +) + + +class MirrorError(RuntimeError): + """Raised when a collision or failed CKAN action makes a mirror unsafe.""" def _slug(text: str, fallback: str) -> str: - """CKAN dataset names must be lowercase slugs, 2-100 chars.""" - s = _SLUG_RE.sub("-", (text or "").lower()).strip("-") - s = re.sub(r"-{2,}", "-", s) - if len(s) < 2: - s = fallback - return s[:100] - - -def _load_index(site: str) -> dict[str, Any]: - """Read the onboarded index.json for a site (local path).""" - path = _PROJECT_ROOT / "data" / "ckan_onboard" / site / "index.json" - if not path.exists(): - raise FileNotFoundError( - f"No onboarded index at {path}. Run scripts/onboard_ckan.py for {site!r} first." - ) - return json.loads(path.read_text()) + """Return a CKAN-safe slug while retaining the legacy helper API.""" + slug = _SLUG_RE.sub("-", (text or "").lower()).strip("-") + slug = re.sub(r"-{2,}", "-", slug) + if len(slug) < 2: + slug = fallback + return slug[:100] + + +def _load_index(site: str, index_path: str | Path | None = None) -> dict[str, Any]: + """Load qsv output when present; mirroring itself does not require it.""" + candidates = ( + [Path(index_path)] + if index_path + else [ + _PROJECT_ROOT / "data" / "ckan_onboard" / site / "index.json", + _PROJECT_ROOT / "ckan_onboard" / site / "index.json", + ] + ) + for path in candidates: + if path.exists(): + return json.loads(path.read_text(encoding="utf-8")) + return {} def _column_dictionary(resource: dict[str, Any]) -> list[dict[str, Any]]: - """Compact, publishable dictionary for one resource's columns.""" - columns = [] - for col in resource.get("columns", []) or []: - stats = col.get("stats", {}) or {} + """Return publishable qsv column metadata with secrets scrubbed.""" + columns: list[dict[str, Any]] = [] + for column in resource.get("columns", []) or []: + stats = column.get("stats", {}) or {} columns.append( { - "name": col.get("name", ""), - "label": col.get("qsv_label") or col.get("ckan_info", {}).get("label", ""), - "description": _scrub_secrets(col.get("qsv_description", "")), - "type": col.get("qsv_type") or col.get("ckan_type", ""), + "name": column.get("name", ""), + "label": column.get("qsv_label") or column.get("ckan_info", {}).get("label", ""), + "description": _scrub_secrets(column.get("qsv_description", "")), + "type": column.get("qsv_type") or column.get("ckan_type", ""), "stats": { - k: stats.get(k) - for k in ("type", "min", "max", "mean", "q2_median", "stddev", - "nullcount", "cardinality") - if stats.get(k) not in (None, "") + key: stats.get(key) + for key in ( + "type", + "min", + "max", + "mean", + "q2_median", + "stddev", + "nullcount", + "cardinality", + ) + if stats.get(key) not in (None, "") }, - "top_values": (col.get("top_values") or [])[:10], + "top_values": (column.get("top_values") or [])[:10], } ) return columns -async def _ensure_org(ckan: CKANClient, org_slug: str, org_title: str) -> None: - """Create the organization if it does not already exist.""" - existing = await ckan.action("organization_show", {"id": org_slug}) - if existing: - return - await ckan.action("organization_create", {"name": org_slug, "title": org_title}) - logger.info("Created organization", org=org_slug) - - -async def _upsert_dataset( - ckan: CKANClient, name: str, payload: dict[str, Any] +def _qsv_by_resource(index: dict[str, Any]) -> dict[str, dict[str, Any]]: + return { + str(resource["resource_id"]): resource + for dataset in index.get("datasets", []) or [] + for resource in dataset.get("resources", []) or [] + if resource.get("resource_id") + } + + +def _project(source: dict[str, Any], fields: Iterable[str]) -> dict[str, Any]: + return { + field: source[field] for field in fields if field in source and source[field] is not None + } + + +def _provenance_extras( + extras: list[dict[str, Any]] | None, + *, + site_id: str, + source_url: str, + source_id: str, +) -> list[dict[str, str]]: + """Preserve source extras and add stable mirror identity fields.""" + values = { + str(item.get("key")): str(item.get("value", "")) + for item in (extras or []) + if item.get("key") + } + values.update( + { + _PROVENANCE_KEYS["portal"]: site_id, + _PROVENANCE_KEYS["url"]: source_url, + _PROVENANCE_KEYS["id"]: source_id, + } + ) + return [{"key": key, "value": value} for key, value in values.items()] + + +def _group_refs(groups: list[dict[str, Any]] | None) -> list[dict[str, str]]: + refs: list[dict[str, str]] = [] + for group in groups or []: + name = group.get("name") + if name: + ref = {"name": str(name)} + if group.get("capacity"): + ref["capacity"] = str(group["capacity"]) + refs.append(ref) + return refs + + +async def _all_packages(client: CKANClient) -> list[dict[str, Any]]: + packages: list[dict[str, Any]] = [] + start = 0 + while True: + page = await client.action( + "package_search", {"q": "*:*", "rows": 1000, "start": start, "sort": "name asc"} + ) + batch = page.get("results", []) if isinstance(page, dict) else [] + packages.extend(batch) + start += len(batch) + if not batch or start >= int(page.get("count", 0)): + return packages + + +async def _source_snapshot( + client: CKANClient, *, organization: str | None = None, limit: int | None = None +) -> dict[str, list[dict[str, Any]]]: + organization_rows = await client.action( + "organization_list", {"all_fields": True, "include_dataset_count": True} + ) + organizations: list[dict[str, Any]] = [] + for row in organization_rows if isinstance(organization_rows, list) else []: + if organization and row.get("name") != organization: + continue + full = await client.action( + "organization_show", + { + "id": row["id"], + "include_datasets": False, + "include_users": False, + "include_groups": True, + "include_tags": True, + }, + ) + organizations.append(full or row) + + group_rows = await client.action("group_list", {"all_fields": True}) + groups = list(group_rows) if isinstance(group_rows, list) else [] + + packages = await _all_packages(client) + if organization: + org_ids = {str(item.get("id")) for item in organizations} + packages = [ + package + for package in packages + if package.get("organization", {}).get("name") == organization + or str(package.get("owner_org")) in org_ids + ] + if limit is not None: + packages = packages[:limit] + used_group_names = { + str(group.get("name")) + for package in packages + for group in package.get("groups", []) or [] + if group.get("name") + } + groups = [group for group in groups if group.get("name") in used_group_names] + return {"organizations": organizations, "groups": groups, "packages": packages} + + +async def _target_snapshot(client: CKANClient) -> dict[str, list[dict[str, Any]]]: + organizations = await client.action("organization_list", {"all_fields": True}) + groups = await client.action("group_list", {"all_fields": True}) + packages = await _all_packages(client) + return { + "organizations": list(organizations) if isinstance(organizations, list) else [], + "groups": list(groups) if isinstance(groups, list) else [], + "packages": packages, + } + + +def _assert_no_collisions( + source: dict[str, list[dict[str, Any]]], + target: dict[str, list[dict[str, Any]]], + *, + store_name: str, + store_id: str, +) -> None: + collisions: list[str] = [] + + def check( + kind: str, + incoming: list[dict[str, Any]], + existing: list[dict[str, Any]], + *, + allow_shared_name: bool = False, + ) -> None: + by_name = {str(row.get("name")): row for row in existing if row.get("name")} + by_id = {str(row.get("id")): row for row in existing if row.get("id")} + for row in incoming: + name, row_id = str(row.get("name", "")), str(row.get("id", "")) + if not allow_shared_name and name in by_name and str(by_name[name].get("id")) != row_id: + collisions.append(f"{kind} name {name!r} already has a different UUID") + if row_id in by_id and str(by_id[row_id].get("name")) != name: + collisions.append(f"{kind} UUID {row_id!r} already has a different name") + + root_existing = next( + (row for row in target["organizations"] if row.get("name") == store_name), None + ) + if root_existing and str(root_existing.get("id")) != store_id: + collisions.append(f"source-store organization {store_name!r} already exists") + check("organization", source["organizations"], target["organizations"]) + # CKAN group/category names are site-wide. Reuse an existing category with + # the same name instead of rewriting it with another portal's UUID. + check("group", source["groups"], target["groups"], allow_shared_name=True) + check("dataset", source["packages"], target["packages"]) + + target_resources = [ + resource + for package in target["packages"] + for resource in package.get("resources", []) or [] + ] + source_package_by_resource = { + str(resource.get("id")): str(package.get("id")) + for package in source["packages"] + for resource in package.get("resources", []) or [] + if resource.get("id") + } + target_by_id = { + str(resource.get("id")): resource for resource in target_resources if resource.get("id") + } + for resource_id in sorted(set(source_package_by_resource) & target_by_id.keys()): + if ( + str(target_by_id[resource_id].get("package_id")) + != source_package_by_resource[resource_id] + ): + collisions.append(f"resource UUID {resource_id!r} belongs to another target dataset") + + if collisions: + detail = "\n - ".join(collisions[:50]) + raise MirrorError(f"Mirror preflight found {len(collisions)} collision(s):\n - {detail}") + + +async def _write_action( + client: CKANClient, action: str, payload: dict[str, Any], *, label: str ) -> dict[str, Any]: - """Create the dataset, or update it in place if the slug already exists.""" - existing = await ckan.action("package_show", {"id": name}) - if existing: - payload["id"] = existing["id"] - return await ckan.action("package_update", payload) - return await ckan.action("package_create", payload) - - -async def populate( - ckan_url: str, api_key: str, site: str, org_slug: str, org_title: str, dry_run: bool -) -> dict[str, int]: - index = _load_index(site) - datasets = index.get("datasets", []) or [] - logger.info("Loaded onboarded index", site=site, datasets=len(datasets)) - - counts = {"datasets": 0, "resources": 0, "skipped": 0} - - if dry_run: - for ds in datasets: - resources = [ - r for r in ds.get("resources", []) if r.get("status") in ("ok", None) - ] - print( - f" would publish: {_slug(ds.get('dataset_title', ''), ds.get('dataset_id', 'x'))} " - f"({len(resources)} resources, " - f"{sum(len(r.get('columns') or []) for r in resources)} columns)" - ) - counts["datasets"] += 1 - counts["resources"] += len(resources) - return counts + result = await client.action(action, payload) + if not result: + raise MirrorError(f"{action} failed for {label}") + return result + + +def _resource_payload( + resource: dict[str, Any], + *, + package_id: str, + site_id: str, + source_url: str, + qsv: dict[str, Any] | None, +) -> dict[str, Any]: + payload = _project(resource, _RESOURCE_FIELDS) + payload["package_id"] = package_id + payload.update( + { + _PROVENANCE_KEYS["portal"]: site_id, + _PROVENANCE_KEYS["url"]: source_url, + _PROVENANCE_KEYS["id"]: str(resource.get("id", "")), + } + ) + # Preserve extension-owned spatial metadata when the destination has the + # same extension, but never claim that a DataStore table was copied. + for key, value in resource.items(): + if key.startswith("dataspatial_") and value is not None: + payload[key] = value + if qsv: + columns = _column_dictionary(qsv) + payload.update( + { + "qsv_description": _scrub_secrets(qsv.get("qsv_description", "")), + "qsv_tags": json.dumps(qsv.get("qsv_tags") or []), + "data_dictionary": json.dumps(columns), + "column_count": len(columns), + "row_count": qsv.get("row_count", 0), + "qsv_onboarded_at": qsv.get("onboarded_at", ""), + } + ) + return payload + + +async def mirror_catalog( + *, + source: CKANClient, + target: CKANClient, + site_id: str, + site_title: str, + source_url: str, + apply: bool = False, + organization: str | None = None, + qsv_index: dict[str, Any] | None = None, + limit: int | None = None, +) -> dict[str, Any]: + """Plan or apply one complete, API-level CKAN metadata mirror.""" + source_data = await _source_snapshot(source, organization=organization, limit=limit) + target_data = await _target_snapshot(target) + store_name = _slug(site_id, "source-store") + store_id = str(uuid.uuid5(uuid.NAMESPACE_URL, source_url.rstrip("/"))) + _assert_no_collisions(source_data, target_data, store_name=store_name, store_id=store_id) + + qsv_resources = _qsv_by_resource(qsv_index or {}) + resources = [ + resource + for package in source_data["packages"] + for resource in package.get("resources", []) or [] + ] + summary: dict[str, Any] = { + "mode": "apply" if apply else "dry-run", + "site": site_id, + "source_url": source_url, + "store_organization": store_name, + "organizations": len(source_data["organizations"]), + "groups": len(source_data["groups"]), + "datasets": len(source_data["packages"]), + "resources": len(resources), + "qsv_resources": sum( + 1 for resource in resources if str(resource.get("id")) in qsv_resources + ), + "created": 0, + "updated": 0, + } + if not apply: + return summary + + target_org_ids = {str(row.get("id")) for row in target_data["organizations"]} + target_group_ids = {str(row.get("id")) for row in target_data["groups"]} + target_group_names = {str(row.get("name")) for row in target_data["groups"]} + target_package_ids = {str(row.get("id")) for row in target_data["packages"]} + target_resource_ids = { + str(resource.get("id")) + for package in target_data["packages"] + for resource in package.get("resources", []) or [] + } + + store_payload = { + "id": store_id, + "name": store_name, + "title": site_title, + "description": f"Datasets mirrored from {source_url.rstrip('/')}", + "extras": _provenance_extras( + [], site_id=site_id, source_url=source_url, source_id=store_id + ), + } + store_exists = store_id in target_org_ids + await _write_action( + target, + "organization_patch" if store_exists else "organization_create", + store_payload, + label=store_name, + ) + summary["updated" if store_exists else "created"] += 1 + + source_org_names = {str(row.get("name")) for row in source_data["organizations"]} + for organization_row in source_data["organizations"]: + payload = _project(organization_row, _ORGANIZATION_FIELDS) + source_id = str(organization_row.get("id", "")) + payload["extras"] = _provenance_extras( + organization_row.get("extras"), + site_id=site_id, + source_url=source_url, + source_id=source_id, + ) + exists = source_id in target_org_ids + await _write_action( + target, + "organization_patch" if exists else "organization_create", + payload, + label=str(organization_row.get("name")), + ) + summary["updated" if exists else "created"] += 1 + + # Apply hierarchy only after every source organization exists. CKAN's + # no_loops validator resolves both child and parent rows; combining a + # caller-supplied child UUID with groups during create trips that validator. + for organization_row in source_data["organizations"]: + parents = _group_refs(organization_row.get("groups")) + if not any(parent["name"] in source_org_names for parent in parents): + parents.append({"name": store_name, "capacity": "parent"}) + await _write_action( + target, + "organization_patch", + {"id": str(organization_row.get("id", "")), "groups": parents}, + label=f"{organization_row.get('name')} hierarchy", + ) - ckan = CKANClient(ckan_url=ckan_url, api_key=api_key) - try: - await _ensure_org(ckan, org_slug, org_title) - - for ds in datasets: - resources = [ - r for r in ds.get("resources", []) if r.get("status") in ("ok", None) - ] - if not resources: - counts["skipped"] += 1 - continue - - name = _slug(ds.get("dataset_title", ""), ds.get("dataset_id", "dataset")) - tags = [{"name": _slug(t, "tag")} for t in (ds.get("tags") or []) if t][:20] - - pkg = await _upsert_dataset( - ckan, - name, - { - "name": name, - "title": ds.get("dataset_title", name), - "notes": _scrub_secrets(ds.get("dataset_description", "")), - "owner_org": org_slug, - "tags": tags, - "extras": [ - {"key": "source_dataset_id", "value": str(ds.get("dataset_id", ""))}, - {"key": "onboarded_from", "value": site}, - ], - }, + for group in source_data["groups"]: + payload = _project(group, _GROUP_FIELDS) + source_id = str(group.get("id", "")) + if str(group.get("name")) in target_group_names and source_id not in target_group_ids: + continue + payload["extras"] = _provenance_extras( + group.get("extras"), site_id=site_id, source_url=source_url, source_id=source_id + ) + exists = source_id in target_group_ids + await _write_action( + target, + "group_patch" if exists else "group_create", + payload, + label=str(group.get("name")), + ) + summary["updated" if exists else "created"] += 1 + + known_org_ids = {str(row.get("id")) for row in source_data["organizations"]} + for package in source_data["packages"]: + payload = _project(package, _PACKAGE_FIELDS) + package_id = str(package.get("id", "")) + owner_org = str(package.get("owner_org") or "") + payload["owner_org"] = owner_org if owner_org in known_org_ids else store_id + payload["tags"] = [ + { + "name": str(tag["name"]), + **({"vocabulary_id": tag["vocabulary_id"]} if tag.get("vocabulary_id") else {}), + } + for tag in package.get("tags", []) or [] + if tag.get("name") + ] + payload["groups"] = _group_refs(package.get("groups")) + payload["extras"] = _provenance_extras( + package.get("extras"), + site_id=site_id, + source_url=source_url, + source_id=package_id, + ) + if not owner_org: + payload["extras"].append({"key": "mirror_source_owner_org", "value": ""}) + exists = package_id in target_package_ids + await _write_action( + target, + "package_patch" if exists else "package_create", + payload, + label=str(package.get("name")), + ) + summary["updated" if exists else "created"] += 1 + + for resource in package.get("resources", []) or []: + resource_id = str(resource.get("id", "")) + resource_payload = _resource_payload( + resource, + package_id=package_id, + site_id=site_id, + source_url=source_url, + qsv=qsv_resources.get(resource_id), ) - if not pkg: - logger.warning("Dataset upsert failed", dataset=name) - counts["skipped"] += 1 - continue - counts["datasets"] += 1 - - for res in resources: - columns = _column_dictionary(res) - await ckan.action( - "resource_create", - { - "package_id": pkg["id"], - "name": res.get("resource_name", "resource"), - "description": _scrub_secrets(res.get("qsv_description", "")), - "format": res.get("format", "CSV"), - "data_dictionary": json.dumps(columns), - "column_count": len(columns), - "row_count": res.get("row_count", 0), - }, - ) - counts["resources"] += 1 - - return counts - finally: - await ckan.close() + resource_exists = resource_id in target_resource_ids + await _write_action( + target, + "resource_patch" if resource_exists else "resource_create", + resource_payload, + label=resource_id, + ) + summary["updated" if resource_exists else "created"] += 1 + return summary -def main() -> None: - ap = argparse.ArgumentParser(description="Publish onboarded dictionaries into CKAN") - ap.add_argument("--ckan-url", default="http://localhost:5001") - ap.add_argument("--api-key", default="") - ap.add_argument("--site", default="wprdc") - ap.add_argument("--org", default="city-of-pittsburgh", help="organization slug") - ap.add_argument("--org-title", default="City of Pittsburgh") - ap.add_argument( - "--dry-run", action="store_true", help="report what would be published, write nothing" - ) - args = ap.parse_args() - if not args.dry_run and not args.api_key: - ap.error("--api-key is required unless --dry-run") +async def _main(args: argparse.Namespace) -> dict[str, Any]: + site = get_site(args.site) + if not site: + raise MirrorError(f"Unknown site {args.site!r}; add it to ckan_sites.json first") + if site.get("portal_type", "ckan") != "ckan": + raise MirrorError("Exact organization mirroring currently requires a CKAN source portal") - counts = asyncio.run( - populate( - args.ckan_url, args.api_key, args.site, args.org, args.org_title, args.dry_run + api_key = os.environ.get(args.api_key_env, "") + if args.api_key_file: + api_key = Path(args.api_key_file).read_text(encoding="utf-8").strip() + if args.apply and not api_key: + raise MirrorError( + f"Set {args.api_key_env} or --api-key-file when using --apply; a sysadmin token is required" ) - ) - print( - f"\nFair Store populate ({'dry run' if args.dry_run else 'done'}): " - f"{counts['datasets']} datasets, {counts['resources']} resources, " - f"{counts['skipped']} skipped." - ) + + source = CKANClient(ckan_url=site["url"]) + target = CKANClient(ckan_url=args.target_url, api_key=api_key or None) + try: + return await mirror_catalog( + source=source, + target=target, + site_id=args.site, + site_title=site["name"], + source_url=site["url"], + apply=args.apply, + organization=args.organization, + qsv_index=_load_index(args.site, args.index_path), + limit=args.limit, + ) + finally: + await source.close() + await target.close() + + +def main() -> None: + parser = argparse.ArgumentParser(description="Mirror CKAN metadata and qsv results") + parser.add_argument("--site", default="wprdc", help="source site ID from ckan_sites.json") + parser.add_argument("--target-url", default=os.environ.get("CKAN_URL", "http://localhost:5001")) + parser.add_argument("--organization", help="optionally mirror only one source organization") + parser.add_argument("--index-path", help="override the qsv index.json path") + parser.add_argument("--limit", type=int, help="limit datasets for a smoke test") + parser.add_argument("--apply", action="store_true", help="write after collision preflight") + parser.add_argument("--api-key-env", default="CKAN_API_KEY") + parser.add_argument("--api-key-file", help="read the target sysadmin token from a file") + args = parser.parse_args() + try: + summary = asyncio.run(_main(args)) + except MirrorError as exc: + parser.exit(2, f"error: {exc}\n") + print(json.dumps(summary, indent=2, sort_keys=True)) if __name__ == "__main__": diff --git a/tests/unit/test_fairstore_populate.py b/tests/unit/test_fairstore_populate.py index 3b61e9b..ed6b25c 100644 --- a/tests/unit/test_fairstore_populate.py +++ b/tests/unit/test_fairstore_populate.py @@ -3,6 +3,8 @@ import importlib.util from pathlib import Path +import pytest + _spec = importlib.util.spec_from_file_location( "populate_fairstore", Path(__file__).resolve().parent.parent.parent / "scripts" / "populate_fairstore.py", @@ -12,6 +14,7 @@ _slug = populate_fairstore._slug _column_dictionary = populate_fairstore._column_dictionary +MirrorError = populate_fairstore.MirrorError class TestSlug: @@ -65,3 +68,202 @@ def test_falls_back_to_ckan_label(self) -> None: col = _column_dictionary(res)[0] assert col["label"] == "Y label" assert col["type"] == "text" + + +class FakeCKAN: + def __init__(self, responses): + self.responses = responses + self.calls = [] + + async def action(self, name, payload=None): + self.calls.append((name, payload or {})) + response = self.responses.get(name, {}) + return response(payload or {}) if callable(response) else response + + +def _source_client(): + organization = { + "id": "11111111-1111-4111-8111-111111111111", + "name": "source-publisher", + "title": "Source Publisher", + "description": "Original publisher description", + "extras": [{"key": "original", "value": "yes"}], + "groups": [], + } + package = { + "id": "22222222-2222-4222-8222-222222222222", + "name": "source-dataset", + "title": "Source Dataset", + "notes": "Original notes", + "owner_org": organization["id"], + "tags": [{"name": "source tag"}], + "groups": [], + "extras": [{"key": "frequency", "value": "monthly"}], + "resources": [ + { + "id": "33333333-3333-4333-8333-333333333333", + "name": "Source CSV", + "url": "https://source.example/data.csv", + "description": "Original resource description", + "format": "CSV", + "datastore_active": True, + } + ], + } + return ( + FakeCKAN( + { + "organization_list": [organization], + "organization_show": organization, + "group_list": [], + "package_search": {"count": 1, "results": [package]}, + } + ), + organization, + package, + ) + + +@pytest.mark.asyncio +async def test_mirror_preserves_source_identity_and_adds_hierarchy_and_qsv(): + source, organization, package = _source_client() + target = FakeCKAN( + { + "organization_list": [], + "group_list": [], + "package_search": {"count": 0, "results": []}, + "organization_create": lambda payload: payload, + "organization_patch": lambda payload: payload, + "package_create": lambda payload: payload, + "resource_create": lambda payload: payload, + } + ) + resource_id = package["resources"][0]["id"] + qsv_index = { + "datasets": [ + { + "resources": [ + { + "resource_id": resource_id, + "qsv_description": "Generated profile", + "qsv_tags": ["finance"], + "row_count": 9, + "columns": [{"name": "amount", "qsv_type": "Float"}], + } + ] + } + ] + } + + summary = await populate_fairstore.mirror_catalog( + source=source, + target=target, + site_id="source-store", + site_title="Source Store", + source_url="https://source.example", + apply=True, + qsv_index=qsv_index, + ) + + assert summary["datasets"] == 1 + assert summary["qsv_resources"] == 1 + child = next( + payload + for action, payload in target.calls + if action == "organization_create" and payload.get("id") == organization["id"] + ) + assert child["name"] == organization["name"] + hierarchy = next( + payload + for action, payload in target.calls + if action == "organization_patch" + and payload.get("id") == organization["id"] + and "groups" in payload + ) + assert hierarchy["groups"] == [{"name": "source-store", "capacity": "parent"}] + + dataset = next(payload for action, payload in target.calls if action == "package_create") + assert dataset["id"] == package["id"] + assert dataset["name"] == package["name"] + assert dataset["notes"] == "Original notes" + assert dataset["tags"] == [{"name": "source tag"}] + + resource = next(payload for action, payload in target.calls if action == "resource_create") + assert resource["id"] == resource_id + assert resource["url"] == "https://source.example/data.csv" + assert resource["description"] == "Original resource description" + assert resource["qsv_description"] == "Generated profile" + assert resource["row_count"] == 9 + assert "datastore_active" not in resource + + +def test_preflight_rejects_dataset_name_collision_but_reuses_group_category(): + source = { + "organizations": [], + "groups": [{"id": "source-group", "name": "education"}], + "packages": [{"id": "source-package", "name": "shared-name", "resources": []}], + } + target = { + "organizations": [], + "groups": [{"id": "target-group", "name": "education"}], + "packages": [{"id": "target-package", "name": "shared-name", "resources": []}], + } + with pytest.raises(MirrorError, match="dataset name 'shared-name'"): + populate_fairstore._assert_no_collisions( + source, + target, + store_name="source-store", + store_id="aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa", + ) + + +@pytest.mark.asyncio +async def test_existing_source_ids_are_patched_instead_of_duplicated(): + source, organization, package = _source_client() + root_id = str( + populate_fairstore.uuid.uuid5( + populate_fairstore.uuid.NAMESPACE_URL, "https://source.example" + ) + ) + target = FakeCKAN( + { + "organization_list": [ + {"id": root_id, "name": "source-store"}, + {"id": organization["id"], "name": organization["name"]}, + ], + "group_list": [], + "package_search": { + "count": 1, + "results": [ + { + "id": package["id"], + "name": package["name"], + "resources": [ + {**resource, "package_id": package["id"]} + for resource in package["resources"] + ], + } + ], + }, + "organization_patch": lambda payload: payload, + "package_patch": lambda payload: payload, + "resource_patch": lambda payload: payload, + } + ) + + await populate_fairstore.mirror_catalog( + source=source, + target=target, + site_id="source-store", + site_title="Source Store", + source_url="https://source.example", + apply=True, + ) + + actions = [action for action, _ in target.calls] + assert "organization_create" not in actions + assert "package_create" not in actions + assert "resource_create" not in actions + assert actions.count("organization_patch") == 3 + assert actions.count("package_patch") == 1 + assert actions.count("resource_patch") == 1