From 2edd6f22fb08f0933e494a71a561d67701d4c68a Mon Sep 17 00:00:00 2001 From: Eric Larson Date: Sat, 3 Oct 2026 00:15:26 -0500 Subject: [PATCH 1/3] Show indexed sources and verify nightly harvest completion --- .github/workflows/ogm-nightly-sync.yml | 7 +- backend/app/api/v1/endpoint_modules/ogm.py | 74 +++++++---------- .../app/services/ogm_harvest/index_status.py | 35 ++++++++ backend/scripts/trigger_ogm_nightly_sync.py | 14 ++++ backend/scripts/wait_ogm_harvest.py | 49 +++++++++++ backend/static/brand.css | 2 +- backend/templates/ogm_repo_dashboard.html | 82 ++++--------------- .../tests/api/v1/test_ogm_public_endpoints.py | 82 +++++++++---------- .../scripts/test_trigger_ogm_nightly_sync.py | 43 ++++++++++ .../tests/scripts/test_wait_ogm_harvest.py | 67 +++++++++++++++ .../tests/services/test_ogm_index_status.py | 39 +++++++++ docs/deployment.md | 16 ++++ 12 files changed, 355 insertions(+), 155 deletions(-) create mode 100644 backend/app/services/ogm_harvest/index_status.py create mode 100644 backend/scripts/wait_ogm_harvest.py create mode 100644 backend/tests/scripts/test_trigger_ogm_nightly_sync.py create mode 100644 backend/tests/scripts/test_wait_ogm_harvest.py create mode 100644 backend/tests/services/test_ogm_index_status.py diff --git a/.github/workflows/ogm-nightly-sync.yml b/.github/workflows/ogm-nightly-sync.yml index 511ec90..9a23435 100644 --- a/.github/workflows/ogm-nightly-sync.yml +++ b/.github/workflows/ogm-nightly-sync.yml @@ -14,8 +14,9 @@ concurrency: jobs: trigger-nightly-sync: - name: Trigger Production OGM Sync + name: Harvest and Reindex Production runs-on: ubuntu-latest + timeout-minutes: 360 steps: - name: Start SSH agent @@ -31,7 +32,7 @@ jobs: mkdir -p "$HOME/.ssh" ssh-keyscan -p "$SSH_PORT" -H "$SSH_HOST" >> "$HOME/.ssh/known_hosts" - - name: Refresh OGM repos and enqueue nightly harvest + - name: Refresh repositories, harvest, and verify indexing env: SSH_HOST: ${{ secrets.OGM_KAMAL_SSH_HOST }} SSH_PORT: ${{ secrets.OGM_KAMAL_SSH_PORT || '22' }} @@ -62,5 +63,5 @@ jobs: fi echo "Using container: $container" - docker exec "$container" python /app/backend/scripts/trigger_ogm_nightly_sync.py + docker exec "$container" python /app/backend/scripts/trigger_ogm_nightly_sync.py --wait ' diff --git a/backend/app/api/v1/endpoint_modules/ogm.py b/backend/app/api/v1/endpoint_modules/ogm.py index af5c9d3..c72b167 100644 --- a/backend/app/api/v1/endpoint_modules/ogm.py +++ b/backend/app/api/v1/endpoint_modules/ogm.py @@ -2,14 +2,16 @@ from html import escape from pathlib import Path from typing import Any, Optional +from urllib.parse import quote, urlencode -from fastapi import APIRouter, Query, Request +from fastapi import APIRouter, HTTPException, Query, Request from fastapi.responses import HTMLResponse from fastapi.templating import Jinja2Templates from app.api.errors import PUBLIC_ERROR_RESPONSES from app.api.schemas import OGMHarvestFailuresResponse, OGMRepoSummariesResponse from app.api.v1.utils import create_response +from app.services.ogm_harvest.index_status import get_index_status from app.services.ogm_harvest.repository import OGMHarvestRepository router = APIRouter() @@ -48,56 +50,43 @@ async def list_public_ogm_repos(): response_class=HTMLResponse, ) async def ogm_repo_dashboard(request: Request): - repos = await ogm_repo.list_public_repo_summaries() - + try: + index_status = await get_index_status() + except Exception as exc: + raise HTTPException( + status_code=503, detail="Search index status is unavailable. Please try again later." + ) from exc + catalog = {r["ogm_repo_name"]: r for r in await ogm_repo.list_public_repo_summaries()} dashboard_repos = [] - counts = await ogm_repo.get_public_dashboard_counts() - repos_with_aardvark = 0 enabled_repos = 0 - never_harvested = 0 - - for repo in repos: - unpublished_count = int(repo.get("unpublished_record_count") or 0) - has_aardvark = bool(repo.get("ogm_has_aardvark")) + for name, count in sorted(index_status["repo_counts"].items()): + if count <= 0: + continue + repo = catalog.get(name, {"ogm_repo_name": name}) enabled = bool(repo.get("ogm_enabled")) and not repo.get("ogm_archived", False) - last_harvest_completed = repo.get("last_crawl_completed_at") - api_hidden_breakdown = [] - if unpublished_count: - api_hidden_breakdown.append({"count": unpublished_count, "label": "unpublished"}) - if repo.get("other_active_source_count"): - api_hidden_breakdown.append( - {"count": repo["other_active_source_count"], "label": "also in active sources"} - ) - - repos_with_aardvark += int(has_aardvark) - enabled_repos += int(enabled) - never_harvested += int(not bool(last_harvest_completed)) - + scheduled = enabled and repo.get("ogm_watch_mode") in { + "nightly", + "weekly", + "scheduled", + "both", + } + enabled_repos += int(scheduled) dashboard_repos.append( { **repo, + "ogm_github_url": repo.get("ogm_github_url") + or f"https://github.com/OpenGeoMetadata/{quote(name, safe='')}", + "ogm_search_url": "/api/v1/search?" + urlencode({"ogm_repo": name}), "display_last_commit_at": _format_timestamp(repo.get("last_commit_at")), - "display_last_harvest_at": _format_timestamp(last_harvest_completed), - "display_last_harvest_started_at": _format_timestamp( - repo.get("last_crawl_started_at") - ), - "source_status": "Archived" - if repo.get("ogm_archived") - else ("Active" if enabled else "Disabled"), - "api_hidden_breakdown": api_hidden_breakdown, - "aardvark_status": {True: "present", False: "missing", None: "unknown"}.get( - repo.get("ogm_has_aardvark"), "unknown" - ), + "display_last_harvest_at": _format_timestamp(repo.get("last_crawl_completed_at")), + "source_status": "Nightly" if scheduled else "Not scheduled", + "indexed_record_count": count, } ) - summary = { "repo_count": len(dashboard_repos), "enabled_repo_count": enabled_repos, - "repos_with_aardvark_count": repos_with_aardvark, - "never_harvested_count": never_harvested, - "harvested_record_count": counts["active_record_count"], - "available_record_count": counts["published_record_count"], + "indexed_record_count": index_status["record_count"], } if templates is None: @@ -110,9 +99,7 @@ async def ogm_repo_dashboard(request: Request): f'Search API' f"{escape(str(repo.get('display_last_commit_at') or '-'))}" f"{escape(str(repo.get('display_last_harvest_at') or '-'))}" - f"{repo['aardvark_status']}" - f"{int(repo.get('harvested_record_count') or 0)}" - f"{int(repo.get('available_record_count') or 0)}" + f"{repo['indexed_record_count']}" "" ) for repo in dashboard_repos @@ -126,8 +113,7 @@ async def ogm_repo_dashboard(request: Request): "

Templates are unavailable, showing a minimal fallback view.

" "" "" - "" - "" + "" f"{rows}
RepositoryLast commitLast harvestAardvarkLast-seen source recordsPublished with source tag (database)Indexed records
" ) ) diff --git a/backend/app/services/ogm_harvest/index_status.py b/backend/app/services/ogm_harvest/index_status.py new file mode 100644 index 0000000..685cd8f --- /dev/null +++ b/backend/app/services/ogm_harvest/index_status.py @@ -0,0 +1,35 @@ +"""Read repository membership directly from the live search index.""" + +import os + +from app.elasticsearch.client import es + + +async def get_index_status() -> dict: + counts = {} + after = None + total = 0 + while True: + composite = { + "size": 100, + "sources": [{"repo": {"terms": {"field": "ogm_repo.keyword"}}}], + } + if after is not None: + composite["after"] = after + response = await es.search( + index=os.getenv("ELASTICSEARCH_INDEX", "opengeometadata_api"), + size=0, + track_total_hits=True, + allow_partial_search_results=False, + aggregations={"repos": {"composite": composite}}, + ) + if response.get("timed_out") or response.get("_shards", {}).get("failed", 0): + raise RuntimeError("Search index returned incomplete repository counts") + total = response["hits"]["total"]["value"] + page = response["aggregations"]["repos"] + for bucket in page["buckets"]: + counts[bucket["key"]["repo"]] = bucket["doc_count"] + after = page.get("after_key") + if not page["buckets"] or after is None: + break + return {"record_count": total, "repo_counts": counts} diff --git a/backend/scripts/trigger_ogm_nightly_sync.py b/backend/scripts/trigger_ogm_nightly_sync.py index ac45a38..d7254a9 100755 --- a/backend/scripts/trigger_ogm_nightly_sync.py +++ b/backend/scripts/trigger_ogm_nightly_sync.py @@ -60,7 +60,12 @@ def main() -> None: help="Refresh the repo catalog but do not enqueue ogm_harvest_all", ) parser.add_argument("--dry-run", action="store_true", help="Do not write or enqueue anything") + parser.add_argument( + "--wait", action="store_true", help="Wait for all harvest and indexing jobs to succeed" + ) args = parser.parse_args() + if args.wait and (args.dry_run or args.skip_harvest or args.limit is not None): + parser.error("--wait requires a complete, non-dry-run harvest") token = args.github_token or os.getenv("GITHUB_TOKEN") database_url = os.getenv("DATABASE_URL") @@ -88,6 +93,8 @@ def main() -> None: try: has_aardvark = repo_has_metadata_aardvark(args.org, name, default_branch, token) except Exception as exc: + if args.wait: + raise RuntimeError(f"Repository discovery failed for {name}") from exc has_aardvark = False repo.setdefault("notes", str(exc)) @@ -105,6 +112,12 @@ def main() -> None: if not args.dry_run and not args.skip_harvest: task = ogm_harvest_all.delay(trigger="nightly") harvest_task_id = task.id + print(f"Harvest batch queued: {harvest_task_id}", flush=True) + if args.wait: + from app.tasks.worker import celery_app + from scripts.wait_ogm_harvest import wait_for_harvest + + wait_for_harvest(task, celery_app.AsyncResult) print( json.dumps( @@ -116,6 +129,7 @@ def main() -> None: "upserted": 0 if args.dry_run else upserted, "harvest_enqueued": bool(harvest_task_id), "harvest_task_id": harvest_task_id, + "harvest_completed": bool(args.wait and harvest_task_id), "dry_run": bool(args.dry_run), }, indent=2, diff --git a/backend/scripts/wait_ogm_harvest.py b/backend/scripts/wait_ogm_harvest.py new file mode 100644 index 0000000..f7a5f07 --- /dev/null +++ b/backend/scripts/wait_ogm_harvest.py @@ -0,0 +1,49 @@ +"""Wait for every repository's harvest and search synchronization to finish.""" + +import time + + +def wait_for_harvest(task, result_factory, timeout=18000): + deadline = time.monotonic() + timeout + + def result(item): + remaining = deadline - time.monotonic() + if remaining <= 0: + raise TimeoutError("Nightly harvest deadline exceeded") + return item.get(timeout=remaining, propagate=True) + + batch = result(task) + ids = batch.get("task_ids", []) + names = batch.get("repo_names", []) + if ( + not ids + or len(ids) != batch.get("enqueued") + or len(ids) != len(names) + or len(set(ids)) != len(ids) + or len(set(names)) != len(names) + ): + raise RuntimeError("Incomplete or empty harvest batch") + failures = [] + for name, task_id in zip(names, ids): + try: + outcome = result(result_factory(task_id)) + stats = outcome.get("stats", {}) + index = stats.get("search_index", {}) + if ( + not outcome.get("ogm_run_id") + or outcome.get("ogm_repo_name") != name + or stats.get("errors") != 0 + or index.get("errors") != 0 + or not {"processed", "indexed", "deleted"}.issubset(index) + ): + raise RuntimeError("Harvest or indexing did not complete successfully") + print( + f"Completed {name}: imported={stats.get('imported', 0)} " + f"indexed={index['indexed']} deleted={index['deleted']}", + flush=True, + ) + except Exception as exc: + failures.append(f"{name}: {exc}") + if failures: + raise RuntimeError("Nightly sync failed: " + "; ".join(failures)) + return batch diff --git a/backend/static/brand.css b/backend/static/brand.css index 7497231..dcc27e3 100644 --- a/backend/static/brand.css +++ b/backend/static/brand.css @@ -1264,7 +1264,7 @@ section#footer-app code:hover { .monitor-summary { display: grid; - grid-template-columns: repeat(6, minmax(0, 1fr)); + grid-template-columns: repeat(3, minmax(0, 1fr)); gap: 0; margin-top: calc(var(--docs-unit) * 1.4); overflow: hidden; diff --git a/backend/templates/ogm_repo_dashboard.html b/backend/templates/ogm_repo_dashboard.html index 5f07442..c614871 100644 --- a/backend/templates/ogm_repo_dashboard.html +++ b/backend/templates/ogm_repo_dashboard.html @@ -6,7 +6,7 @@ {{ title }} - +
@@ -35,40 +35,25 @@

OpenGeoMetadata Harvest Monitor

Repository health, harvest freshness, and API availability in one place.

- This page tracks every repository discovered in the OpenGeoMetadata GitHub organization, - whether it exposes a metadata-aardvark/ directory, when it last changed on GitHub, - when this app last harvested it, and how many records were last seen in each source. Archived and disabled - sources retain their harvest history. Publication counts come from the database; - this page does not verify the live search index. + Repositories with records in the live search index. Each nightly sync fetches the latest + metadata, imports it, and updates the index. Use Search API to browse a repository’s records.

Generated {{ generated_at }}. Times shown in UTC.

-

Tracked Repositories

+

Indexed Repositories

{{ summary.repo_count }}

-

With Aardvark Metadata

-

{{ summary.repos_with_aardvark_count }}

+

Searchable Records

+

{{ "{:,}".format(summary.indexed_record_count) }}

-

Harvesting Enabled

+

Scheduled Nightly

{{ summary.enabled_repo_count }}

-
-

Unique Records in Active Sources

-

{{ "{:,}".format(summary.harvested_record_count) }}

-
-
-

Published Records (Database)

-

{{ "{:,}".format(summary.available_record_count) }}

-
-
-

Never Harvested Here

-

{{ summary.never_harvested_count }}

-
@@ -76,8 +61,8 @@

Repository health, harvest freshness, and API availability in one place.

Repositories

Per-repository counts can overlap and should not be added together. - “Also in active sources” describes overlapping source records, not a synchronization failure.

-

JSON source: /api/v1/ogm/repos

+ Counts come directly from the live search index.

+

Full repository catalog: /api/v1/ogm/repos

@@ -88,10 +73,8 @@

Repositories

Repository Last Commit Last Harvest - Aardvark Last Harvest Result - Last-Seen Source Records - Published with Source Tag (Database) + Indexed Records @@ -102,13 +85,7 @@

Repositories

{{ repo.ogm_repo_name }}

{{ repo.ogm_repo_full_name or repo.ogm_repo_name }}

Search API

-

{{ repo.source_status }}

-

- mode: {{ repo.ogm_watch_mode or "manual" }} - {% if repo.last_commit_sha %} - · sha {{ repo.last_commit_sha[:12] }} - {% endif %} -

+

{{ repo.source_status }}

@@ -117,22 +94,10 @@

Repositories

{% if repo.display_last_harvest_at %}
{{ repo.display_last_harvest_at }}
- {% if repo.display_last_harvest_started_at and repo.display_last_harvest_started_at != repo.display_last_harvest_at %} -

started {{ repo.display_last_harvest_started_at }}

- {% endif %} {% else %} Never harvested {% endif %} - - {% if repo.ogm_has_aardvark is sameas true %} - Present - {% elif repo.ogm_has_aardvark is sameas false %} - Missing - {% else %} - Unknown - {% endif %} - {% set crawl_status = (repo.last_crawl_status or "unknown")|lower %} {% if crawl_status == "success" %} @@ -148,24 +113,10 @@

Repositories

{{ "{:,}".format(repo.harvested_failure_count) }} recent import errors

{% endif %} - -
{{ "{:,}".format(repo.harvested_record_count) }}
- {% if repo.harvested_success_count %} -

last run imported {{ "{:,}".format(repo.harvested_success_count) }}

- {% endif %} - - -
{{ "{:,}".format(repo.available_record_count) }}
- {% if repo.api_hidden_breakdown %} -

- {% for item in repo.api_hidden_breakdown %} - {{ "{:,}".format(item.count) }} {{ item.label }}{% if not loop.last %} · {% endif %} - {% endfor %} -

- - {% endif %} - + {{ "{:,}".format(repo.indexed_record_count) }} + {% else %} + No repositories currently have indexed records. {% endfor %} @@ -191,8 +142,9 @@

Monitor Links

Notes

    -
  • Harvested record counts reflect current OGM repo state tracked in ogm_resource_state.
  • -
  • API visibility counts include published records exposed through this API, including records suppressed from GeoBlacklight search results; hidden counts identify unpublished or unsynced records.
  • +
  • Only repositories represented in the search index appear above.
  • +
  • The searchable total includes records without a repository tag.
  • +
  • Nightly runs refresh metadata and synchronize each repository’s search records.
diff --git a/backend/tests/api/v1/test_ogm_public_endpoints.py b/backend/tests/api/v1/test_ogm_public_endpoints.py index ff65b6d..80ad0a6 100644 --- a/backend/tests/api/v1/test_ogm_public_endpoints.py +++ b/backend/tests/api/v1/test_ogm_public_endpoints.py @@ -68,9 +68,9 @@ def test_public_ogm_repo_dashboard_renders_html_monitor(): ogm.ogm_repo, "list_public_repo_summaries", AsyncMock(return_value=sample_repos) ), patch.object( - ogm.ogm_repo, - "get_public_dashboard_counts", - AsyncMock(return_value={"active_record_count": 904, "published_record_count": 1234}), + ogm, + "get_index_status", + AsyncMock(return_value={"record_count": 1234, "repo_counts": {"edu.utexas": 899}}), ), ): response = client.get("/api/v1/ogm/repos/dashboard") @@ -81,7 +81,9 @@ def test_public_ogm_repo_dashboard_renders_html_monitor(): assert "edu.utexas" in response.text assert "OpenGeoMetadata/edu.utexas" in response.text assert "/api/v1/ogm/repos" in response.text - assert "3 unpublished" in response.text + assert "899" in response.text + assert "1,234" in response.text + assert "Never Harvested Here" not in response.text assert 'href="https://github.com/OpenGeoMetadata/edu.utexas"' in response.text assert 'href="/api/v1/search?ogm_repo=edu.utexas"' in response.text @@ -169,57 +171,53 @@ def test_public_ogm_failures_endpoint_can_limit_to_hard_failures_only(): ) -def test_dashboard_distinguishes_retired_sources_and_uses_global_counts(): +def test_dashboard_only_shows_indexed_sources_even_when_database_disagrees(): repos = [ - { - "ogm_repo_name": "edu.umn", - "ogm_enabled": True, - "ogm_archived": True, - "last_crawl_status": "success", - "harvested_record_count": 6332, - "available_record_count": 0, - "unpublished_record_count": 441, - "other_active_source_count": 5891, - }, - { - "ogm_repo_name": "gov.usgs", - "ogm_enabled": False, - "harvested_record_count": 177789, - "available_record_count": 0, - "other_active_source_count": 177789, - }, + {"ogm_repo_name": "retired", "ogm_enabled": False, "available_record_count": 100}, + {"ogm_repo_name": "empty", "ogm_enabled": True, "available_record_count": 10}, { "ogm_repo_name": "geobtaa", "ogm_enabled": True, - "harvested_record_count": 40703, - "available_record_count": 40703, + "ogm_watch_mode": "nightly", + "available_record_count": 1, }, ] with ( patch.object(ogm.ogm_repo, "list_public_repo_summaries", AsyncMock(return_value=repos)), patch.object( - ogm.ogm_repo, - "get_public_dashboard_counts", + ogm, + "get_index_status", AsyncMock( - return_value={"active_record_count": 315135, "published_record_count": 318408} + return_value={ + "record_count": 318408, + "repo_counts": {"geobtaa": 40703, "uncataloged": 2}, + } ), ), ): response = client.get("/api/v1/ogm/repos/dashboard") assert response.status_code == 200 - html = response.text - for expected in ( - "Archived", - "Disabled", - "Active", - "315,135", - "318,408", - "5,891 also in active sources", - "177,789 also in active sources", - "Last Harvest Result", - "Published Records (Database)", + for expected in ("318,408", "40,703", "geobtaa", "uncataloged", "Scheduled Nightly"): + assert expected in response.text + for absent in ("retired", "Never Harvested Here", "Database", "also in active sources"): + assert absent not in response.text + assert 'href="https://github.com/OpenGeoMetadata/uncataloged"' in response.text + + +def test_dashboard_reports_index_unavailable_instead_of_misleading_empty_counts(): + with patch.object(ogm, "get_index_status", AsyncMock(side_effect=RuntimeError("offline"))): + response = client.get("/api/v1/ogm/repos/dashboard") + assert response.status_code == 503 + assert "unavailable" in response.json()["detail"] + + +def test_dashboard_empty_index(): + with ( + patch.object(ogm.ogm_repo, "list_public_repo_summaries", AsyncMock(return_value=[])), + patch.object( + ogm, "get_index_status", AsyncMock(return_value={"record_count": 0, "repo_counts": {}}) + ), ): - assert expected in html - assert "not yet synced" not in html - assert "hidden from API" not in html - assert "does not verify the live search index" in html + response = client.get("/api/v1/ogm/repos/dashboard") + assert response.status_code == 200 + assert "No repositories currently have indexed records" in response.text diff --git a/backend/tests/scripts/test_trigger_ogm_nightly_sync.py b/backend/tests/scripts/test_trigger_ogm_nightly_sync.py new file mode 100644 index 0000000..b624e16 --- /dev/null +++ b/backend/tests/scripts/test_trigger_ogm_nightly_sync.py @@ -0,0 +1,43 @@ +from unittest.mock import Mock, patch + +import pytest + +from scripts import trigger_ogm_nightly_sync as trigger + + +def test_wait_checks_children_after_enqueue(monkeypatch): + monkeypatch.setenv("DATABASE_URL", "postgresql://unused/test") + monkeypatch.setattr("sys.argv", ["trigger", "--wait"]) + task = Mock(id="batch") + with ( + patch.object(trigger, "list_org_repos", return_value=[]), + patch.object(trigger, "upsert_rows", return_value=(0, 0)), + patch.object(trigger.ogm_harvest_all, "delay", return_value=task) as enqueue, + patch("scripts.wait_ogm_harvest.wait_for_harvest") as wait, + ): + trigger.main() + enqueue.assert_called_once_with(trigger="nightly") + assert wait.call_args.args[0] is task + + +def test_discovery_error_does_not_disable_sources_or_enqueue(monkeypatch): + monkeypatch.setenv("DATABASE_URL", "postgresql://unused/test") + monkeypatch.setattr("sys.argv", ["trigger", "--wait"]) + with ( + patch.object(trigger, "list_org_repos", return_value=[{"name": "source"}]), + patch.object(trigger, "repo_has_metadata_aardvark", side_effect=RuntimeError("offline")), + patch.object(trigger, "upsert_rows") as upsert, + patch.object(trigger.ogm_harvest_all, "delay") as enqueue, + ): + with pytest.raises(RuntimeError, match="Repository discovery failed"): + trigger.main() + upsert.assert_not_called() + enqueue.assert_not_called() + + +@pytest.mark.parametrize("flag", ["--dry-run", "--skip-harvest", "--limit=1"]) +def test_wait_requires_a_complete_batch(monkeypatch, flag): + monkeypatch.setattr("sys.argv", ["trigger", "--wait", flag]) + with pytest.raises(SystemExit) as exc: + trigger.main() + assert exc.value.code == 2 diff --git a/backend/tests/scripts/test_wait_ogm_harvest.py b/backend/tests/scripts/test_wait_ogm_harvest.py new file mode 100644 index 0000000..c5a5513 --- /dev/null +++ b/backend/tests/scripts/test_wait_ogm_harvest.py @@ -0,0 +1,67 @@ +from unittest.mock import Mock, patch + +import pytest + +from scripts.wait_ogm_harvest import wait_for_harvest + + +def successful(name="a"): + return { + "ogm_run_id": 1, + "ogm_repo_name": name, + "stats": { + "errors": 0, + "imported": 5, + "search_index": {"processed": 5, "indexed": 5, "deleted": 0, "errors": 0}, + }, + } + + +def batch(): + return Mock( + get=Mock(return_value={"enqueued": 2, "repo_names": ["a", "b"], "task_ids": ["1", "2"]}) + ) + + +def test_waits_for_all_children_and_index_results(): + factory = Mock( + side_effect=[ + Mock(get=Mock(return_value=successful("a"))), + Mock(get=Mock(return_value=successful("b"))), + ] + ) + assert wait_for_harvest(batch(), factory)["enqueued"] == 2 + assert factory.call_count == 2 + + +@pytest.mark.parametrize("failure", ["import", "index", "missing_index", "skipped", "raised"]) +def test_failed_child_fails_batch_but_remaining_children_are_awaited(failure): + outcome = successful() + if failure == "import": + outcome["stats"]["errors"] = 1 + if failure == "index": + outcome["stats"]["search_index"]["errors"] = 1 + if failure == "missing_index": + del outcome["stats"]["search_index"] + if failure == "skipped": + outcome = {"status": "skipped"} + first = ( + Mock(get=Mock(side_effect=RuntimeError("worker failed"))) + if failure == "raised" + else Mock(get=Mock(return_value=outcome)) + ) + factory = Mock(side_effect=[first, Mock(get=Mock(return_value=successful("b")))]) + with pytest.raises(RuntimeError, match="Nightly sync failed"): + wait_for_harvest(batch(), factory) + assert factory.call_count == 2 + + +def test_empty_batch_cannot_report_success(): + with pytest.raises(RuntimeError, match="Incomplete or empty"): + wait_for_harvest(Mock(get=Mock(return_value={"enqueued": 0})), Mock()) + + +def test_wait_deadline(): + with patch("scripts.wait_ogm_harvest.time.monotonic", side_effect=[0, 100]): + with pytest.raises(TimeoutError): + wait_for_harvest(batch(), Mock(), timeout=10) diff --git a/backend/tests/services/test_ogm_index_status.py b/backend/tests/services/test_ogm_index_status.py new file mode 100644 index 0000000..4495379 --- /dev/null +++ b/backend/tests/services/test_ogm_index_status.py @@ -0,0 +1,39 @@ +from unittest.mock import AsyncMock, patch + +import pytest + +from app.services.ogm_harvest import index_status + + +@pytest.mark.asyncio +async def test_index_counts_follow_all_composite_pages_and_exact_total(): + def page(buckets, after=None): + return { + "hits": {"total": {"value": 30}}, + "aggregations": { + "repos": { + "buckets": [ + {"key": {"repo": name}, "doc_count": count} for name, count in buckets + ], + "after_key": after, + } + }, + } + + search = AsyncMock(side_effect=[page([("a", 10)], {"repo": "a"}), page([("b", 12)])]) + with patch.object(index_status.es, "search", search): + result = await index_status.get_index_status() + assert result == {"record_count": 30, "repo_counts": {"a": 10, "b": 12}} + assert search.call_args_list[0].kwargs["track_total_hits"] is True + assert search.call_args_list[0].kwargs["allow_partial_search_results"] is False + assert search.call_args_list[1].kwargs["aggregations"]["repos"]["composite"]["after"] == { + "repo": "a" + } + + +@pytest.mark.asyncio +@pytest.mark.parametrize("response", [{"timed_out": True}, {"_shards": {"failed": 1}}]) +async def test_partial_counts_are_rejected(response): + with patch.object(index_status.es, "search", AsyncMock(return_value=response)): + with pytest.raises(RuntimeError, match="incomplete"): + await index_status.get_index_status() diff --git a/docs/deployment.md b/docs/deployment.md index be8e8ac..b637e15 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -803,3 +803,19 @@ PostgreSQL records. Database rows removed entirely are outside this helper's sco Dashboard harvested totals count repository memberships, so shared IDs can appear more than once. Its published total covers repository-tagged database records; search also includes published records without a repository tag. + +### Nightly harvest and indexing verification + +The Nightly OGM Sync GitHub Actions workflow runs daily at 07:15 UTC (02:15 CDT, +01:15 CST); GitHub may delay scheduled starts. It refreshes the repository catalog, +harvests every enabled scheduled source, and synchronizes each source into the +live search index. The workflow uses `trigger_ogm_nightly_sync.py --wait` and only +succeeds after every child harvest reports zero import errors and a completed +search synchronization with zero indexing errors. Failures and a five-hour wait +timeout fail the workflow. The existing concurrency group prevents overlapping +workflow runs while their harvests are being monitored. + +Deploy the updated trigger and wait scripts before running the updated workflow. +Container-based nightly cron remains disabled to avoid a second scheduler. +The repository dashboard lists sources with records in Elasticsearch, using live +index counts rather than database publication counts. From f331a67c556dbf3f9a4a609b093726867e2ee174 Mon Sep 17 00:00:00 2001 From: Eric Larson Date: Sat, 3 Oct 2026 00:16:41 -0500 Subject: [PATCH 2/3] Use a concise dashboard heading --- backend/templates/ogm_repo_dashboard.html | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/backend/templates/ogm_repo_dashboard.html b/backend/templates/ogm_repo_dashboard.html index c614871..8085e19 100644 --- a/backend/templates/ogm_repo_dashboard.html +++ b/backend/templates/ogm_repo_dashboard.html @@ -33,7 +33,7 @@

OpenGeoMetadata Harvest Monitor

-

Repository health, harvest freshness, and API availability in one place.

+

Repositories in the search index.

Repositories with records in the live search index. Each nightly sync fetches the latest metadata, imports it, and updates the index. Use Search API to browse a repository’s records. From f0c275a37d963348146f2a980288823c16394d32 Mon Sep 17 00:00:00 2001 From: Eric Larson Date: Sat, 3 Oct 2026 00:18:57 -0500 Subject: [PATCH 3/3] Keep the nightly SSH connection alive during long harvests --- .github/workflows/ogm-nightly-sync.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ogm-nightly-sync.yml b/.github/workflows/ogm-nightly-sync.yml index 9a23435..b76ef98 100644 --- a/.github/workflows/ogm-nightly-sync.yml +++ b/.github/workflows/ogm-nightly-sync.yml @@ -38,7 +38,7 @@ jobs: SSH_PORT: ${{ secrets.OGM_KAMAL_SSH_PORT || '22' }} SSH_USER: ${{ secrets.OGM_KAMAL_SSH_USER }} run: | - ssh -o BatchMode=yes -p "$SSH_PORT" "$SSH_USER@$SSH_HOST" ' + ssh -o BatchMode=yes -o ServerAliveInterval=30 -o ServerAliveCountMax=6 -p "$SSH_PORT" "$SSH_USER@$SSH_HOST" ' set -eu container="$(docker ps \