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
9 changes: 5 additions & 4 deletions .github/workflows/ogm-nightly-sync.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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 \
Expand All @@ -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
'
74 changes: 30 additions & 44 deletions backend/app/api/v1/endpoint_modules/ogm.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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:
Expand All @@ -110,9 +99,7 @@ async def ogm_repo_dashboard(request: Request):
f'<a href="{escape(repo["ogm_search_url"], quote=True)}">Search API</a></td>'
f"<td>{escape(str(repo.get('display_last_commit_at') or '-'))}</td>"
f"<td>{escape(str(repo.get('display_last_harvest_at') or '-'))}</td>"
f"<td>{repo['aardvark_status']}</td>"
f"<td>{int(repo.get('harvested_record_count') or 0)}</td>"
f"<td>{int(repo.get('available_record_count') or 0)}</td>"
f"<td>{repo['indexed_record_count']}</td>"
"</tr>"
)
for repo in dashboard_repos
Expand All @@ -126,8 +113,7 @@ async def ogm_repo_dashboard(request: Request):
"<p>Templates are unavailable, showing a minimal fallback view.</p>"
"<table><thead><tr>"
"<th>Repository</th><th>Last commit</th><th>Last harvest</th>"
"<th>Aardvark</th><th>Last-seen source records</th>"
"<th>Published with source tag (database)</th>"
"<th>Indexed records</th>"
f"</tr></thead><tbody>{rows}</tbody></table></body></html>"
)
)
Expand Down
35 changes: 35 additions & 0 deletions backend/app/services/ogm_harvest/index_status.py
Original file line number Diff line number Diff line change
@@ -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}
14 changes: 14 additions & 0 deletions backend/scripts/trigger_ogm_nightly_sync.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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))

Expand All @@ -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(
Expand All @@ -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,
Expand Down
49 changes: 49 additions & 0 deletions backend/scripts/wait_ogm_harvest.py
Original file line number Diff line number Diff line change
@@ -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
2 changes: 1 addition & 1 deletion backend/static/brand.css
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading
Loading