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 --- API.md | 13 ++++++ README.md | 6 ++- alembic/versions/20260411_0002_tracked_actors.py | 39 ++++++++++++++++ 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 ++++++++++++++-- tests/test_contributions_api.py | 21 ++++++++- tests/test_ingestion.py | 30 +++++++++++- tests/test_migrations.py | 3 ++ tests/test_polling_api.py | 22 ++++++++- 13 files changed, 258 insertions(+), 15 deletions(-) create mode 100644 alembic/versions/20260411_0002_tracked_actors.py diff --git a/API.md b/API.md index cca4b8d..da38662 100644 --- a/API.md +++ b/API.md @@ -43,6 +43,7 @@ Dates: Contribution ranges: - Contribution read endpoints return zero-filled days, so clients do not need to patch missing dates. +- Public contribution reads also register the actor for background indexing. A first lookup may therefore return an empty graph while the scheduler catches up. Repository names: @@ -114,6 +115,9 @@ Response `200 OK`: "actor": "~your-user", "from": "2026-03-01", "to": "2026-04-15", + "is_indexed": false, + "last_polled_at": null, + "indexing_state": "pending", "days": [ { "date": "2026-03-01", "count": 0, "score": 0.0 }, { "date": "2026-03-02", "count": 3, "score": 2.5 } @@ -126,6 +130,9 @@ Response fields: - `actor` string: canonical actor after alias resolution - `from` string: inclusive start date - `to` string: inclusive end date +- `is_indexed` boolean: whether the service has already indexed activity for this actor +- `last_polled_at` string or `null`: most recent successful poll time, if any +- `indexing_state` string: one of `pending`, `indexed`, or `error` - `days` array: - `date` string `YYYY-MM-DD` - `count` integer contribution count for the day @@ -174,6 +181,9 @@ Response `200 OK`: "actor": "~your-user", "from": "2026-03-01", "to": "2026-04-15", + "is_indexed": true, + "last_polled_at": "2026-04-11T18:05:00Z", + "indexing_state": "indexed", "total_events": 126, "total_score": 116.75, "active_days": 14, @@ -187,6 +197,9 @@ Response fields: - `actor` string - `from` string - `to` string +- `is_indexed` boolean +- `last_polled_at` string or `null` +- `indexing_state` string - `total_events` integer - `total_score` float - `active_days` integer diff --git a/README.md b/README.md index 97c9b43..bb4ac7c 100644 --- a/README.md +++ b/README.md @@ -165,7 +165,7 @@ Example response: } ``` -Scheduled polling only runs when `ENABLE_SCHEDULER=true` and uses `DEFAULT_ACTOR`. +Scheduled polling only runs when `ENABLE_SCHEDULER=true`. The scheduler seeds `DEFAULT_ACTOR` as an initial known actor, and public contribution reads register additional actors for later background polling. For `git.sr.ht`, tracked repositories are configured via `GIT_TRACKED_REPOSITORIES`. Entries may be either: @@ -205,6 +205,9 @@ Example response: "actor": "~your-user", "from": "2026-01-01", "to": "2026-03-30", + "is_indexed": true, + "last_polled_at": "2026-04-11T18:05:00Z", + "indexing_state": "indexed", "days": [ {"date": "2026-03-28", "count": 3, "score": 3.5}, {"date": "2026-03-29", "count": 0, "score": 0.0}, @@ -313,6 +316,7 @@ The SourceHut-specific assumptions are isolated to the service modules: - `git.sr.ht` polling is limited to repositories listed in `GIT_TRACKED_REPOSITORIES` - scheduled polling runs in-process, so it is not a distributed scheduler +- newly requested actors are indexed asynchronously, so the first public read may be empty until a scheduler or manual poll runs - alias management is config-driven; there is no alias CRUD API yet - current deployment model is trusted-operator V1, not a public multi-tenant service diff --git a/alembic/versions/20260411_0002_tracked_actors.py b/alembic/versions/20260411_0002_tracked_actors.py new file mode 100644 index 0000000..e3f0e31 --- /dev/null +++ b/alembic/versions/20260411_0002_tracked_actors.py @@ -0,0 +1,39 @@ +"""tracked actors for lazy indexing""" + +from __future__ import annotations + +from alembic import op +import sqlalchemy as sa +from sqlalchemy import inspect + + +revision = "20260411_0002" +down_revision = "20260409_0001" +branch_labels = None +depends_on = None + + +def _table_names() -> set[str]: + return set(inspect(op.get_bind()).get_table_names()) + + +def upgrade() -> None: + if "tracked_actors" in _table_names(): + return + + op.create_table( + "tracked_actors", + sa.Column("id", sa.Integer(), primary_key=True), + sa.Column("actor", sa.String(length=255), nullable=False), + sa.Column("is_active", sa.Boolean(), nullable=False, server_default=sa.true()), + sa.Column("last_requested_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("last_polled_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("last_poll_status", sa.String(length=32), nullable=True), + sa.Column("last_poll_error", sa.Text(), nullable=True), + sa.UniqueConstraint("actor", name="uq_tracked_actor_actor"), + ) + + +def downgrade() -> None: + if "tracked_actors" in _table_names(): + op.drop_table("tracked_actors") 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]: diff --git a/tests/test_contributions_api.py b/tests/test_contributions_api.py index e221602..22d0af3 100644 --- a/tests/test_contributions_api.py +++ b/tests/test_contributions_api.py @@ -1,9 +1,10 @@ from datetime import UTC, datetime from fastapi.testclient import TestClient +from sqlalchemy import select from srht_contrib.main import create_app -from srht_contrib.models import ContributionEvent +from srht_contrib.models import ContributionEvent, TrackedActor def test_read_only_contribution_routes_are_public_and_write_routes_require_api_key(settings, db_engine, session_factory) -> None: @@ -45,6 +46,8 @@ def test_contributions_api_returns_zero_filled_range(client: TestClient, db_sess response = client.get("/api/contributions/~ccleberg?from=2026-03-28&to=2026-03-30") assert response.status_code == 200 + assert response.json()["is_indexed"] is True + assert response.json()["indexing_state"] == "indexed" assert response.json()["days"] == [ {"date": "2026-03-28", "count": 0, "score": 0.0}, {"date": "2026-03-29", "count": 0, "score": 0.0}, @@ -52,6 +55,20 @@ def test_contributions_api_returns_zero_filled_range(client: TestClient, db_sess ] +def test_public_read_registers_actor_for_lazy_indexing(client: TestClient, db_session) -> None: + response = client.get("/api/contributions/~ccleberg?from=2026-03-28&to=2026-03-30") + + tracked_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg")) + + assert response.status_code == 200 + assert response.json()["is_indexed"] is False + assert response.json()["indexing_state"] == "pending" + assert response.json()["last_polled_at"] is None + assert tracked_actor is not None + assert tracked_actor.is_active is True + assert tracked_actor.last_requested_at is not None + + def test_contribution_stats_api(client: TestClient, db_session) -> None: db_session.add_all( [ @@ -88,6 +105,8 @@ def test_contribution_stats_api(client: TestClient, db_session) -> None: assert response.json()["total_score"] == 1.5 assert response.json()["longest_streak"] == 2 assert response.json()["current_streak"] == 2 + assert response.json()["is_indexed"] is True + assert response.json()["indexing_state"] == "indexed" def test_invalid_date_input_returns_400(client: TestClient) -> None: diff --git a/tests/test_ingestion.py b/tests/test_ingestion.py index 2322b4a..788ec81 100644 --- a/tests/test_ingestion.py +++ b/tests/test_ingestion.py @@ -4,7 +4,7 @@ from sqlalchemy import select from srht_contrib.config import Settings from srht_contrib.jobs.poller import PollerService -from srht_contrib.models import SyncState, TrackedRepository +from srht_contrib.models import SyncState, TrackedActor, TrackedRepository from srht_contrib.schemas import NormalizedEvent from srht_contrib.services.git import GitIngestionService, GitPollResult from srht_contrib.services.todo import TodoIngestionService, TodoPollResult @@ -335,3 +335,31 @@ def test_sync_overlap_reuses_cursor_window_and_suppresses_duplicates(db_session) assert state is not None assert len(todo_service.calls) == 2 assert todo_service.calls[1].isoformat() == "2026-03-30T00:00:00+00:00" + + +def test_scheduled_poll_polls_known_actors_and_seeds_default_actor(db_session) -> None: + event = NormalizedEvent( + service="todo", + event_type="ticket_created", + actor="~known", + repo_name="todo", + resource_id="123", + external_uid="todo:event:known:created:123", + occurred_at=datetime(2026, 3, 30, 10, 0, tzinfo=UTC), + weight=1.0, + raw_payload_json=None, + ) + todo_service = RecordingTodoService(events_by_call=[[], [event]]) + poller = PollerService(todo_service=todo_service, git_service=EmptyGitService()) + + db_session.add(TrackedActor(actor="~known", is_active=True)) + db_session.commit() + + results = poller.poll_tracked_actors(db_session, default_actor="~default") + + tracked_actors = db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all() + + assert results == {"~default": 0, "~known": 1} + assert [actor.actor for actor in tracked_actors] == ["~default", "~known"] + assert all(actor.last_poll_status == "indexed" for actor in tracked_actors) + assert all(actor.last_polled_at is not None for actor in tracked_actors) diff --git a/tests/test_migrations.py b/tests/test_migrations.py index d572af2..505ddb3 100644 --- a/tests/test_migrations.py +++ b/tests/test_migrations.py @@ -22,6 +22,7 @@ def test_alembic_upgrade_creates_schema(tmp_path) -> None: "alembic_version", "contribution_events", "sync_states", + "tracked_actors", "tracked_repositories", ] @@ -115,6 +116,7 @@ def test_alembic_upgrade_adopts_legacy_schema(tmp_path) -> None: assert columns["actor"]["nullable"] is False assert "uq_tracked_repository_service_actor_name" in unique_constraints assert actor == Settings().default_actor + assert "tracked_actors" in inspector.get_table_names() def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) -> None: @@ -128,3 +130,4 @@ def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) -> inspector = inspect(create_engine(database_url)) assert "actor_aliases" in inspector.get_table_names() + assert "tracked_actors" in inspector.get_table_names() diff --git a/tests/test_polling_api.py b/tests/test_polling_api.py index b07b877..75270e0 100644 --- a/tests/test_polling_api.py +++ b/tests/test_polling_api.py @@ -1,10 +1,11 @@ from datetime import UTC, datetime from fastapi.testclient import TestClient +from sqlalchemy import select from srht_contrib.config import Settings from srht_contrib.main import create_app -from srht_contrib.models import ContributionEvent +from srht_contrib.models import ContributionEvent, TrackedActor from srht_contrib.services.srht_client import SourceHutClientError @@ -19,6 +20,16 @@ class InsertingPoller: self.todo_service = service self.git_service = service + def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True): + 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) + if update_last_requested: + tracked_actor.last_requested_at = datetime(2026, 3, 30, 9, 0, tzinfo=UTC) + db.flush() + return tracked_actor + def poll_all(self, db, actor: str) -> int: db.add( ContributionEvent( @@ -43,6 +54,14 @@ class FailingPoller: self.todo_service = service self.git_service = service + def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True): + 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) + db.flush() + return tracked_actor + def poll_all(self, db, actor: str) -> int: raise SourceHutClientError("boom") @@ -58,6 +77,7 @@ def test_manual_poll_uses_same_database_session(settings: Settings, db_engine, s assert poll_response.status_code == 200 assert poll_response.json()["inserted_events"] == 1 assert calendar_response.status_code == 200 + assert calendar_response.json()["is_indexed"] is True assert calendar_response.json()["days"] == [{"date": "2026-03-30", "count": 1, "score": 1.0}] -- cgit v1.2.3