diff --git a/.github/workflows/ogm-nightly-sync.yml b/.github/workflows/ogm-nightly-sync.yml index 511ec90..b76ef98 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,13 +32,13 @@ 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' }} 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 \ @@ -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..8085e19 100644 --- a/backend/templates/ogm_repo_dashboard.html +++ b/backend/templates/ogm_repo_dashboard.html @@ -6,7 +6,7 @@ {{ title }} - +
@@ -33,42 +33,27 @@

OpenGeoMetadata Harvest Monitor

-

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

+

Repositories in the search index.

- 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.