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
6 changes: 5 additions & 1 deletion Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ FROM python:3.11-slim

WORKDIR /app

RUN apt-get update && apt-get install -y \
RUN apt-get update && DEBIAN_FRONTEND=noninteractive apt-get install -y \
gcc \
g++ \
python3-dev \
Expand All @@ -14,7 +14,9 @@ RUN apt-get update && apt-get install -y \
libgdal-dev \
curl \
ca-certificates \
cron \
git \
tzdata \
&& rm -rf /var/lib/apt/lists/*

ENV GDAL_VERSION=3.4.1
Expand All @@ -28,8 +30,10 @@ ENV UV_HTTP_TIMEOUT=300
COPY backend/pyproject.toml backend/uv.lock ./
COPY backend/scripts ./scripts
COPY backend/ ./backend/
COPY config/crontab ./config/crontab

RUN uv pip install -e ./backend --system
RUN chmod +x /app/scripts/start_cron.sh

RUN mkdir -p logs static/maps

Expand Down
2 changes: 2 additions & 0 deletions backend/app/api/v1/strong_params.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@
"include_filters[field][]",
"exclude_filters[field][]",
# Convenience multi-select repo filter (OGM)
"ogm_repo",
"ogm_repo[]",
]

Expand Down Expand Up @@ -72,6 +73,7 @@
"include_filters[field][]",
"exclude_filters[field][]",
# Convenience multi-select repo filter (OGM)
"ogm_repo",
"ogm_repo[]",
]

Expand Down
8 changes: 7 additions & 1 deletion backend/app/elasticsearch/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,13 @@ async def init_elasticsearch():
if "ogm_repo" not in props:
logger.info("Adding missing mapping field: ogm_repo")
await es.indices.put_mapping(
index=index_name, properties={"ogm_repo": {"type": "keyword"}}
index=index_name,
properties={
"ogm_repo": {
"type": "text",
"fields": {"keyword": {"type": "keyword", "ignore_above": 256}},
}
},
)
except Exception as e:
logger.warning(f"Could not ensure mappings for {index_name}: {e}")
Expand Down
94 changes: 79 additions & 15 deletions backend/app/elasticsearch/index.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@

from app.services.language_service import ensure_b1g_language
from db.database import database
from db.models import resources
from db.models import ogm_resource_state, resources

from .client import es
from .suggest import build_suggest_inputs
Expand Down Expand Up @@ -125,6 +125,29 @@ def to_int(v):
return iv if iv is not None else None


def _coerce_ogm_repo_values(value):
"""Normalize OGM repo values for exact faceting/filtering."""
if value in (None, ""):
return []
if isinstance(value, (list, tuple, set)):
raw_values = value
else:
raw_values = [value]

repos = []
seen = set()
for raw in raw_values:
text = str(raw).strip()
if not text:
continue
if text.startswith("ogm_repo:"):
text = text[len("ogm_repo:") :].strip()
if text and text not in seen:
seen.add(text)
repos.append(text)
return repos


def _calculate_time_period_from_year(year_value):
"""Calculate the time period bucket for a given year value.

Expand Down Expand Up @@ -208,7 +231,7 @@ async def index_resources():

await init_elasticsearch()

resource_rows = await database.fetch_all(resources.select())
resource_rows = await fetch_resources_for_index()
processed_resources = await prepare_bulk_data(resource_rows, index_name)

if processed_resources:
Expand All @@ -217,6 +240,42 @@ async def index_resources():
return {"message": "No resources to index"}


async def fetch_resources_for_index():
"""Fetch resource rows and attach current OGM repo memberships."""
resource_rows = await database.fetch_all(resources.select())
ogm_rows = await database.fetch_all(
ogm_resource_state.select().with_only_columns(
ogm_resource_state.c.ogm_resource_id,
ogm_resource_state.c.ogm_repo_name,
ogm_resource_state.c.ogm_missing_since,
)
)

repos_by_resource_id = {}
resources_with_ogm_state = set()
for row in ogm_rows:
row_dict = dict(row)
resource_id = str(row_dict.get("ogm_resource_id") or "").strip()
repo_name = str(row_dict.get("ogm_repo_name") or "").strip()
if not resource_id or not repo_name:
continue
resources_with_ogm_state.add(resource_id)
if row_dict.get("ogm_missing_since") is not None:
continue
repos = repos_by_resource_id.setdefault(resource_id, [])
if repo_name not in repos:
repos.append(repo_name)

indexed_rows = []
for row in resource_rows:
resource = dict(row)
resource_id = str(resource.get("id"))
if resource_id in resources_with_ogm_state:
resource["ogm_repo"] = repos_by_resource_id.get(resource_id, [])
indexed_rows.append(resource)
return indexed_rows


async def prepare_bulk_data(resources, index_name):
"""Prepare resources for indexing (now using individual operations for reliability)."""
processed_resources = []
Expand All @@ -230,6 +289,7 @@ async def prepare_bulk_data(resources, index_name):
async def process_resource(resource_dict):
"""Process a single resource for indexing."""
processed_dict = {}
explicit_ogm_repo_field = "ogm_repo" in resource_dict

date_fields = {"gbl_mdmodified_dt", "b1g_dateAccessioned_s", "b1g_dateRetired_s"}
integer_fields = {"gbl_indexYear_im"}
Expand All @@ -240,7 +300,11 @@ async def process_resource(resource_dict):
}

for key, value in resource_dict.items():
if isinstance(value, (list, tuple)):
if key == "ogm_repo":
repos = _coerce_ogm_repo_values(value)
if repos:
processed_dict[key] = repos
elif isinstance(value, (list, tuple)):
processed_dict[key] = list(value)
elif key in date_fields:
processed_dict[key] = _coerce_date(value)
Expand Down Expand Up @@ -294,27 +358,27 @@ async def process_resource(resource_dict):

ensure_b1g_language(processed_dict)

# Derive OGM repo facet/filter field from admin tags.
explicit_ogm_repo_values = _coerce_ogm_repo_values(processed_dict.get("ogm_repo"))
if explicit_ogm_repo_values:
processed_dict["ogm_repo"] = explicit_ogm_repo_values
else:
processed_dict.pop("ogm_repo", None)

# Derive OGM repo facet/filter field from admin tags when no explicit
# repo state was attached to the index row.
# Source-of-truth tag format stored in Postgres: "ogm_repo:<repo_name>"
tags = processed_dict.get("b1g_adminTags_sm")
if tags:
if tags and not explicit_ogm_repo_field and "ogm_repo" not in processed_dict:
if isinstance(tags, str):
tags_list = [tags]
elif isinstance(tags, list):
tags_list = [str(t) for t in tags if t is not None]
else:
tags_list = [str(tags)]

ogm_repo_values = []
seen = set()
for t in tags_list:
if not isinstance(t, str):
continue
if t.startswith("ogm_repo:"):
repo_name = t[len("ogm_repo:") :].strip()
if repo_name and repo_name not in seen:
seen.add(repo_name)
ogm_repo_values.append(repo_name)
ogm_repo_values = _coerce_ogm_repo_values(
[t for t in tags_list if isinstance(t, str) and t.startswith("ogm_repo:")]
)
if ogm_repo_values:
processed_dict["ogm_repo"] = ogm_repo_values

Expand Down
2 changes: 1 addition & 1 deletion backend/app/elasticsearch/mappings.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@
"type": "text",
"fields": {"keyword": {"type": "keyword", "ignore_above": 8191}},
},
# OpenGeoMetadata repo facet/filter (derived at index-time from b1g_adminTags_sm)
# OpenGeoMetadata repo facet/filter (attached from OGM state at index-time)
"ogm_repo": {
"type": "text",
"fields": {"keyword": {"type": "keyword", "ignore_above": 256}},
Expand Down
6 changes: 4 additions & 2 deletions backend/app/services/search_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -509,8 +509,10 @@ def extract_new_style_filters(self, params: Optional[str]) -> tuple[Dict, Dict]:
list(raw_params.keys())[:10],
)

# Convenience filters (non-bracket style) for common client use cases.
# Example: ogm_repo[]=edu.stanford.purl&ogm_repo[]=edu.umn
# Convenience filters for common client use cases.
# Examples: ogm_repo=edu.unr, ogm_repo[]=edu.stanford.purl&ogm_repo[]=edu.umn
if "ogm_repo" in raw_params:
include_filters.setdefault("ogm_repo", []).extend(raw_params.get("ogm_repo") or [])
if "ogm_repo[]" in raw_params:
include_filters.setdefault("ogm_repo", []).extend(raw_params.get("ogm_repo[]") or [])

Expand Down
70 changes: 70 additions & 0 deletions backend/scripts/ogm_importer.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,66 @@
logger = logging.getLogger(__name__)


def derive_repo_alias(repo_name: str) -> Optional[str]:
parts = [p for p in (repo_name or "").split(".") if p]
if len(parts) >= 2 and parts[0] == "edu":
return parts[1]
if parts:
return parts[0]
return None


def derive_repo_name_from_path(path: str, ogm_path: Optional[str] = None) -> Optional[str]:
path_obj = Path(path)
candidate_parts = list(path_obj.parts)

if ogm_path:
try:
candidate_parts = list(
path_obj.resolve(strict=False)
.relative_to(Path(ogm_path).resolve(strict=False))
.parts
)
except ValueError:
pass

for parts in (candidate_parts, list(path_obj.parts)):
if "metadata-aardvark" not in parts:
continue
index = parts.index("metadata-aardvark")
if index > 0:
return parts[index - 1]

return None


def inject_ogm_repo_tags(record: Dict[str, Any], repo_name: Optional[str]) -> Dict[str, Any]:
if not repo_name:
return record

existing = record.get("b1g_adminTags_sm")
tags: List[str] = []
if isinstance(existing, list):
tags.extend([str(tag).strip() for tag in existing if str(tag).strip()])
elif isinstance(existing, str) and existing.strip():
tags.append(existing.strip())

tags.append(f"ogm_repo:{repo_name}")
if alias := derive_repo_alias(repo_name):
tags.append(f"ogm:{alias}")

deduped: List[str] = []
seen = set()
for tag in tags:
if tag in seen:
continue
seen.add(tag)
deduped.append(tag)

record["b1g_adminTags_sm"] = deduped
return record


class OGMImporter:
"""Imports OpenGeoMetadata Aardvark records into the database."""

Expand Down Expand Up @@ -368,6 +428,16 @@ def import_records(self, limit: Optional[int] = None) -> Dict[str, int]:
stats["skipped"] += 1
continue

repo_name = derive_repo_name_from_path(path, self.ogm_path)
if repo_name:
inject_ogm_repo_tags(prepared_record, repo_name)
else:
logger.warning(
"Unable to derive OGM repo name for record %s from path %s",
prepared_record.get("id"),
path,
)

# Normalize to ensure every model column has a value (or None)
normalized_record = {}
for col in resources.c:
Expand Down
36 changes: 33 additions & 3 deletions backend/scripts/populate_ogm_repos.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import argparse
import json
import os
import sys
from typing import Any, Dict, List, Optional, Tuple
from urllib.parse import urlparse, urlunparse

Expand All @@ -31,6 +32,9 @@
# Keep script self-contained: import the SQLAlchemy Table definitions.
from db.models import ogm_repos

_BAD_GITHUB_TOKEN_WARNING_SHOWN = False
_REJECTED_GITHUB_TOKENS: set[str] = set()


def _sync_database_url(database_url: str) -> str:
# Convert asyncpg URL to sync URL
Expand Down Expand Up @@ -62,14 +66,40 @@ def _github_headers(token: Optional[str]) -> Dict[str, str]:
return headers


def _warn_bad_github_token() -> None:
global _BAD_GITHUB_TOKEN_WARNING_SHOWN
if _BAD_GITHUB_TOKEN_WARNING_SHOWN:
return

print(
"Warning: configured GitHub token was rejected with 401; "
"retrying public GitHub API requests without authentication.",
file=sys.stderr,
)
_BAD_GITHUB_TOKEN_WARNING_SHOWN = True


def _github_get(url: str, token: Optional[str], **kwargs: Any) -> requests.Response:
effective_token = token
if token and token in _REJECTED_GITHUB_TOKENS:
effective_token = None

resp = requests.get(url, headers=_github_headers(effective_token), **kwargs)
if effective_token and resp.status_code == 401:
_REJECTED_GITHUB_TOKENS.add(effective_token)
_warn_bad_github_token()
resp = requests.get(url, headers=_github_headers(None), **kwargs)
return resp


def list_org_repos(org: str, token: Optional[str], per_page: int = 100) -> List[Dict[str, Any]]:
repos: List[Dict[str, Any]] = []
page = 1
while True:
url = f"https://api.github.com/orgs/{org}/repos"
resp = requests.get(
resp = _github_get(
url,
headers=_github_headers(token),
token,
params={"per_page": per_page, "page": page},
timeout=30,
)
Expand All @@ -93,7 +123,7 @@ def repo_has_metadata_aardvark(
params = {}
if default_branch:
params["ref"] = default_branch
resp = requests.get(url, headers=_github_headers(token), params=params, timeout=30)
resp = _github_get(url, token, params=params, timeout=30)
if resp.status_code == 200:
body = resp.json()
return isinstance(body, list)
Expand Down
Loading
Loading