Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 61 additions & 12 deletions scripts/populate_fairstore.py
Original file line number Diff line number Diff line change
Expand Up @@ -210,12 +210,31 @@ async def _all_packages(client: CKANClient) -> list[dict[str, Any]]:
return packages


async def _all_group_rows(client: CKANClient, action_name: str) -> list[dict[str, Any]]:
"""Read every organization or group despite portal-side page-size caps."""
rows: list[dict[str, Any]] = []
seen_ids: set[str] = set()
offset = 0
while True:
result = await client.action(
action_name,
{"all_fields": True, "limit": 1000, "offset": offset},
)
batch = list(result) if isinstance(result, list) else []
if not batch:
return rows
new_rows = [row for row in batch if str(row.get("id")) not in seen_ids]
if not new_rows:
return rows
rows.extend(new_rows)
seen_ids.update(str(row.get("id")) for row in new_rows)
offset += len(batch)


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}
)
organization_rows = await _all_group_rows(client, "organization_list")
organizations: list[dict[str, Any]] = []
for row in organization_rows if isinstance(organization_rows, list) else []:
if organization and row.get("name") != organization:
Expand All @@ -232,8 +251,7 @@ async def _source_snapshot(
)
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 []
groups = await _all_group_rows(client, "group_list")

packages = await _all_packages(client)
if organization:
Expand All @@ -257,8 +275,8 @@ async def _source_snapshot(


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})
organizations = await _all_group_rows(client, "organization_list")
groups = await _all_group_rows(client, "group_list")
packages = await _all_packages(client)
return {
"organizations": list(organizations) if isinstance(organizations, list) else [],
Expand Down Expand Up @@ -297,7 +315,14 @@ def check(
)
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"])
if any(str(row.get("id")) == store_id for row in source["organizations"]):
collisions.append(f"source organization UUID {store_id!r} matches the source-store UUID")
# A prior run may have used the bare site ID for the synthetic root. It is
# renamed before source organizations are written, so exclude it here.
target_organizations = [
row for row in target["organizations"] if str(row.get("id")) != store_id
]
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)
Expand Down Expand Up @@ -332,10 +357,30 @@ def check(
async def _write_action(
client: CKANClient, action: str, payload: dict[str, Any], *, label: str
) -> dict[str, Any]:
result = await client.action(action, payload)
if not result:
raise MirrorError(f"{action} failed for {label}")
return result
show_actions = {
"organization_create": "organization_show",
"group_create": "group_show",
"package_create": "package_show",
"resource_create": "resource_show",
}
for attempt in range(3):
result = await client.action(action, payload)
if result:
return result

# A gateway can time out after CKAN commits the object. Resolve by the
# preserved UUID before retrying so a successful write is not treated
# as a failed duplicate create.
show_action = show_actions.get(action)
if show_action and payload.get("id"):
existing = await client.action(show_action, {"id": payload["id"]})
if existing:
return existing

if attempt < 2:
await asyncio.sleep(2**attempt)

raise MirrorError(f"{action} failed for {label}")


def _resource_payload(
Expand Down Expand Up @@ -391,6 +436,8 @@ async def mirror_catalog(
source_data = await _source_snapshot(source, organization=organization, limit=limit)
target_data = await _target_snapshot(target)
store_name = _slug(site_id, "source-store")
if store_name in {str(row.get("name")) for row in source_data["organizations"]}:
store_name = _slug(f"{store_name}-source", "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)

Expand Down Expand Up @@ -517,6 +564,8 @@ async def mirror_catalog(
source_url=source_url,
source_id=package_id,
)
if not str(payload.get("notes") or "").strip():
payload["notes"] = "No description was provided by the source catalog."
if not owner_org:
payload["extras"].append({"key": "mirror_source_owner_org", "value": ""})
exists = package_id in target_package_ids
Expand Down
87 changes: 87 additions & 0 deletions tests/unit/test_fairstore_populate.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,38 @@ async def action(self, name, payload=None):
return response(payload or {}) if callable(response) else response


@pytest.mark.asyncio
async def test_group_rows_follow_portal_page_caps():
rows = [{"id": str(index), "name": f"publisher-{index}"} for index in range(5)]

def capped_page(payload):
offset = payload.get("offset", 0)
return rows[offset : offset + 2]

client = FakeCKAN({"organization_list": capped_page})

result = await populate_fairstore._all_group_rows(client, "organization_list")

assert result == rows
assert [payload["offset"] for _, payload in client.calls] == [0, 2, 4, 5]


@pytest.mark.asyncio
async def test_write_action_recovers_when_gateway_times_out_after_create():
resource = {"id": "resource-id", "name": "Created resource"}
client = FakeCKAN({"resource_create": {}, "resource_show": resource})

result = await populate_fairstore._write_action(
client,
"resource_create",
resource,
label="resource-id",
)

assert result == resource
assert [action for action, _ in client.calls] == ["resource_create", "resource_show"]


def _source_client():
organization = {
"id": "11111111-1111-4111-8111-111111111111",
Expand Down Expand Up @@ -197,6 +229,61 @@ async def test_mirror_preserves_source_identity_and_adds_hierarchy_and_qsv():
assert "datastore_active" not in resource


@pytest.mark.asyncio
async def test_mirror_supplies_description_when_source_notes_are_blank():
source, _organization, package = _source_client()
package["notes"] = ""
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,
}
)

await populate_fairstore.mirror_catalog(
source=source,
target=target,
site_id="source-store",
site_title="Source Store",
source_url="https://source.example",
apply=True,
)

dataset = next(payload for action, payload in target.calls if action == "package_create")
assert dataset["notes"] == "No description was provided by the source catalog."


@pytest.mark.asyncio
async def test_source_store_slug_does_not_replace_same_named_source_organization():
source, _organization, _package = _source_client()
source_url = "https://source.example"
root_id = str(
populate_fairstore.uuid.uuid5(populate_fairstore.uuid.NAMESPACE_URL, source_url)
)
target = FakeCKAN(
{
"organization_list": [{"id": root_id, "name": "source-publisher"}],
"group_list": [],
"package_search": {"count": 0, "results": []},
}
)

summary = await populate_fairstore.mirror_catalog(
source=source,
target=target,
site_id="source-publisher",
site_title="Source Publisher",
source_url=source_url,
)

assert summary["store_organization"] == "source-publisher-source"


def test_preflight_rejects_dataset_name_collision_but_reuses_group_category():
source = {
"organizations": [],
Expand Down
Loading