From 6149ca9c394df7d316f80db05ad0fc2aea6c550d Mon Sep 17 00:00:00 2001 From: Christian Cleberg Date: Sat, 11 Apr 2026 13:58:01 -0500 Subject: feat: add lazy actor indexing for contribution lookups --- src/srht_contrib/api/routes_contributions.py | 12 +++++- src/srht_contrib/jobs/poller.py | 59 +++++++++++++++++++++++++++- src/srht_contrib/main.py | 2 +- src/srht_contrib/models.py | 15 ++++++- src/srht_contrib/schemas.py | 11 +++++- src/srht_contrib/services/aggregator.py | 40 +++++++++++++++++-- 6 files changed, 128 insertions(+), 11 deletions(-) (limited to 'src') diff --git a/src/srht_contrib/api/routes_contributions.py b/src/srht_contrib/api/routes_contributions.py index c682553..25bc3c2 100644 --- a/src/srht_contrib/api/routes_contributions.py +++ b/src/srht_contrib/api/routes_contributions.py @@ -41,12 +41,16 @@ def get_contributions( year: int | None = Query(default=None, ge=1970, le=3000), from_date: str | None = Query(default=None, alias="from"), to_date: str | None = Query(default=None, alias="to"), + poller: PollerService = Depends(get_poller), db: Session = Depends(get_db), actor_identity_resolver: ActorIdentityResolver = Depends(get_actor_identity_resolver), ) -> ContributionCalendarResponse: start, end = _resolve_range(year, from_date, to_date) canonical_actor = actor_identity_resolver.canonicalize(actor, db=db) - return ContributionAggregator().build_calendar(db, canonical_actor, start, end) + poller.track_actor_request(db, canonical_actor) + response = ContributionAggregator().build_calendar(db, canonical_actor, start, end) + db.commit() + return response @router.get("/{actor}/stats", response_model=ContributionStatsResponse) @@ -55,12 +59,16 @@ def get_contribution_stats( year: int | None = Query(default=None, ge=1970, le=3000), from_date: str | None = Query(default=None, alias="from"), to_date: str | None = Query(default=None, alias="to"), + poller: PollerService = Depends(get_poller), db: Session = Depends(get_db), actor_identity_resolver: ActorIdentityResolver = Depends(get_actor_identity_resolver), ) -> ContributionStatsResponse: start, end = _resolve_range(year, from_date, to_date) canonical_actor = actor_identity_resolver.canonicalize(actor, db=db) - return ContributionAggregator().build_stats(db, canonical_actor, start, end) + poller.track_actor_request(db, canonical_actor) + response = ContributionAggregator().build_stats(db, canonical_actor, start, end) + db.commit() + return response @router.post("/poll", response_model=PollResponse, dependencies=[Depends(require_api_key)]) diff --git a/src/srht_contrib/jobs/poller.py b/src/srht_contrib/jobs/poller.py index c3ed069..910f1ab 100644 --- a/src/srht_contrib/jobs/poller.py +++ b/src/srht_contrib/jobs/poller.py @@ -7,9 +7,10 @@ from sqlalchemy import select from sqlalchemy.exc import IntegrityError from sqlalchemy.orm import Session -from srht_contrib.models import ContributionEvent, SyncState, TrackedRepository +from srht_contrib.models import ContributionEvent, SyncState, TrackedActor, TrackedRepository from srht_contrib.schemas import NormalizedEvent from srht_contrib.services.git import GitIngestionService +from srht_contrib.services.srht_client import SourceHutClientError from srht_contrib.services.todo import TodoIngestionService from srht_contrib.utils.repositories import canonicalize_repository_name @@ -24,6 +25,52 @@ class PollerService: self.git_service = git_service def poll_all(self, db: Session, actor: str) -> int: + self.track_actor_request(db, actor, update_last_requested=False) + try: + inserted = self._poll_actor(db, actor) + except Exception as exc: + db.rollback() + self._update_tracked_actor_poll_state(db, actor, status="error", error=str(exc)) + db.commit() + raise + + self._update_tracked_actor_poll_state(db, actor, status="indexed", error=None) + db.commit() + return inserted + + def poll_tracked_actors(self, db: Session, default_actor: str | None = None) -> dict[str, int]: + if default_actor: + self.track_actor_request(db, default_actor, update_last_requested=False) + db.commit() + + results: dict[str, int] = {} + actors = db.scalars( + select(TrackedActor.actor) + .where(TrackedActor.is_active.is_(True)) + .order_by(TrackedActor.last_requested_at.is_(None), TrackedActor.last_requested_at.desc(), TrackedActor.actor) + ).all() + for actor in actors: + try: + results[actor] = self.poll_all(db, actor) + except SourceHutClientError: + logger.exception("Scheduled poll failed for actor=%s", actor) + except Exception: + logger.exception("Unexpected scheduled poll failure for actor=%s", actor) + return results + + def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor: + tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) + if tracked_actor is None: + tracked_actor = TrackedActor(actor=actor, is_active=True) + db.add(tracked_actor) + + tracked_actor.is_active = True + if update_last_requested: + tracked_actor.last_requested_at = datetime.now(tz=UTC) + db.flush() + return tracked_actor + + def _poll_actor(self, db: Session, actor: str) -> int: inserted = 0 inserted += self._poll_service(db, actor, self.todo_service.service_name, self.todo_service.fetch_recent_events) self._sync_tracked_repositories(db, actor) @@ -38,7 +85,6 @@ class PollerService: repositories=git_repositories, ), ) - db.commit() return inserted def _poll_service(self, db: Session, actor: str, service_name: str, fetcher) -> int: @@ -117,3 +163,12 @@ class PollerService: .order_by(TrackedRepository.repo_name) ).all() return list(rows) + + def _update_tracked_actor_poll_state(self, db: Session, actor: str, status: str, error: str | None) -> None: + tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) + tracked_actor.last_poll_status = status + tracked_actor.last_poll_error = error + if status == "indexed": + tracked_actor.last_polled_at = datetime.now(tz=UTC) + db.add(tracked_actor) + db.flush() diff --git a/src/srht_contrib/main.py b/src/srht_contrib/main.py index 37c3492..9047bcc 100644 --- a/src/srht_contrib/main.py +++ b/src/srht_contrib/main.py @@ -91,7 +91,7 @@ def _scheduled_poll(app: FastAPI) -> None: session_factory: sessionmaker[Session] = app.state.session_factory db = session_factory() try: - poller.poll_all(db, settings.default_actor) + poller.poll_tracked_actors(db, settings.default_actor) finally: db.close() diff --git a/src/srht_contrib/models.py b/src/srht_contrib/models.py index b87f464..2f33482 100644 --- a/src/srht_contrib/models.py +++ b/src/srht_contrib/models.py @@ -2,7 +2,7 @@ from __future__ import annotations from datetime import datetime -from sqlalchemy import JSON, DateTime, Float, Index, Integer, String, Text, UniqueConstraint +from sqlalchemy import JSON, Boolean, DateTime, Float, Index, Integer, String, Text, UniqueConstraint from sqlalchemy.orm import Mapped, mapped_column from srht_contrib.db import Base @@ -58,3 +58,16 @@ class ActorAlias(Base): id: Mapped[int] = mapped_column(Integer, primary_key=True) canonical_actor: Mapped[str] = mapped_column(String(255), nullable=False) alias: Mapped[str] = mapped_column(String(255), nullable=False) + + +class TrackedActor(Base): + __tablename__ = "tracked_actors" + __table_args__ = (UniqueConstraint("actor", name="uq_tracked_actor_actor"),) + + id: Mapped[int] = mapped_column(Integer, primary_key=True) + actor: Mapped[str] = mapped_column(String(255), nullable=False) + is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True) + last_requested_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + last_polled_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + last_poll_status: Mapped[str | None] = mapped_column(String(32), nullable=True) + last_poll_error: Mapped[str | None] = mapped_column(Text, nullable=True) diff --git a/src/srht_contrib/schemas.py b/src/srht_contrib/schemas.py index 2a00996..bebf93d 100644 --- a/src/srht_contrib/schemas.py +++ b/src/srht_contrib/schemas.py @@ -1,6 +1,7 @@ from __future__ import annotations from datetime import date, datetime +from typing import Literal from pydantic import BaseModel, Field, model_validator @@ -23,7 +24,13 @@ class ContributionDay(BaseModel): score: float -class ContributionCalendarResponse(BaseModel): +class ContributionIndexMetadata(BaseModel): + is_indexed: bool + last_polled_at: datetime | None = None + indexing_state: Literal["pending", "indexed", "error"] + + +class ContributionCalendarResponse(ContributionIndexMetadata): actor: str from_date: date = Field(alias="from") to_date: date = Field(alias="to") @@ -32,7 +39,7 @@ class ContributionCalendarResponse(BaseModel): model_config = {"populate_by_name": True} -class ContributionStatsResponse(BaseModel): +class ContributionStatsResponse(ContributionIndexMetadata): actor: str from_date: date = Field(alias="from") to_date: date = Field(alias="to") diff --git a/src/srht_contrib/services/aggregator.py b/src/srht_contrib/services/aggregator.py index 800329a..90ef93b 100644 --- a/src/srht_contrib/services/aggregator.py +++ b/src/srht_contrib/services/aggregator.py @@ -6,8 +6,13 @@ from datetime import date from sqlalchemy import func, select from sqlalchemy.orm import Session -from srht_contrib.models import ContributionEvent -from srht_contrib.schemas import ContributionCalendarResponse, ContributionDay, ContributionStatsResponse +from srht_contrib.models import ContributionEvent, TrackedActor +from srht_contrib.schemas import ( + ContributionCalendarResponse, + ContributionDay, + ContributionIndexMetadata, + ContributionStatsResponse, +) from srht_contrib.utils.dates import date_range, date_to_utc_bounds @@ -30,7 +35,14 @@ class ContributionAggregator: ) for day in date_range(start, end) ] - return ContributionCalendarResponse(actor=actor, from_date=start, to_date=end, days=days) + metadata = self._index_metadata(db, actor) + return ContributionCalendarResponse( + actor=actor, + from_date=start, + to_date=end, + days=days, + **metadata.model_dump(), + ) def build_stats(self, db: Session, actor: str, start: date, end: date) -> ContributionStatsResponse: calendar = self.build_calendar(db, actor, start, end) @@ -47,6 +59,28 @@ class ContributionAggregator: active_days=len(active_days), longest_streak=max(streaks, default=0), current_streak=current_streak, + is_indexed=calendar.is_indexed, + last_polled_at=calendar.last_polled_at, + indexing_state=calendar.indexing_state, + ) + + def _index_metadata(self, db: Session, actor: str) -> ContributionIndexMetadata: + tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) + has_indexed_events = db.scalar(select(ContributionEvent.id).where(ContributionEvent.actor == actor).limit(1)) is not None + last_poll_status = tracked_actor.last_poll_status if tracked_actor is not None else None + is_indexed = has_indexed_events or (tracked_actor is not None and tracked_actor.last_polled_at is not None) + + if last_poll_status == "error": + indexing_state = "error" + elif is_indexed: + indexing_state = "indexed" + else: + indexing_state = "pending" + + return ContributionIndexMetadata( + is_indexed=is_indexed, + last_polled_at=tracked_actor.last_polled_at if tracked_actor is not None else None, + indexing_state=indexing_state, ) def _query_daily_aggregates(self, db: Session, actor: str, start: date, end: date) -> list[DailyAggregate]: -- cgit v1.2.3