krz/hutch-stats

Server-side utility for calculating contributions for sourcehut users. contributions server sourcehut stats

Commit 6149ca9c39

6149ca9c394df7d316f80db05ad0fc2aea6c550d

parent: b6d1348abf

Verified · cmc

cmc <hello@cleberg.net> · 2026-04-11 18:58 UTC

feat: add lazy actor indexing for contribution lookups

Layout: unified · split

API.md +13
@@ -43,6 +43,7 @@ Dates:
43Contribution ranges: 43Contribution ranges:
44 44
45- Contribution read endpoints return zero-filled days, so clients do not need to patch missing dates. 45- Contribution read endpoints return zero-filled days, so clients do not need to patch missing dates.
46- Public contribution reads also register the actor for background indexing. A first lookup may therefore return an empty graph while the scheduler catches up.
46 47
47Repository names: 48Repository names:
48 49
@@ -114,6 +115,9 @@ Response `200 OK`:
114 "actor": "~your-user", 115 "actor": "~your-user",
115 "from": "2026-03-01", 116 "from": "2026-03-01",
116 "to": "2026-04-15", 117 "to": "2026-04-15",
118 "is_indexed": false,
119 "last_polled_at": null,
120 "indexing_state": "pending",
117 "days": [ 121 "days": [
118 { "date": "2026-03-01", "count": 0, "score": 0.0 }, 122 { "date": "2026-03-01", "count": 0, "score": 0.0 },
119 { "date": "2026-03-02", "count": 3, "score": 2.5 } 123 { "date": "2026-03-02", "count": 3, "score": 2.5 }
@@ -126,6 +130,9 @@ Response fields:
126- `actor` string: canonical actor after alias resolution 130- `actor` string: canonical actor after alias resolution
127- `from` string: inclusive start date 131- `from` string: inclusive start date
128- `to` string: inclusive end date 132- `to` string: inclusive end date
133- `is_indexed` boolean: whether the service has already indexed activity for this actor
134- `last_polled_at` string or `null`: most recent successful poll time, if any
135- `indexing_state` string: one of `pending`, `indexed`, or `error`
129- `days` array: 136- `days` array:
130 - `date` string `YYYY-MM-DD` 137 - `date` string `YYYY-MM-DD`
131 - `count` integer contribution count for the day 138 - `count` integer contribution count for the day
@@ -174,6 +181,9 @@ Response `200 OK`:
174 "actor": "~your-user", 181 "actor": "~your-user",
175 "from": "2026-03-01", 182 "from": "2026-03-01",
176 "to": "2026-04-15", 183 "to": "2026-04-15",
184 "is_indexed": true,
185 "last_polled_at": "2026-04-11T18:05:00Z",
186 "indexing_state": "indexed",
177 "total_events": 126, 187 "total_events": 126,
178 "total_score": 116.75, 188 "total_score": 116.75,
179 "active_days": 14, 189 "active_days": 14,
@@ -187,6 +197,9 @@ Response fields:
187- `actor` string 197- `actor` string
188- `from` string 198- `from` string
189- `to` string 199- `to` string
200- `is_indexed` boolean
201- `last_polled_at` string or `null`
202- `indexing_state` string
190- `total_events` integer 203- `total_events` integer
191- `total_score` float 204- `total_score` float
192- `active_days` integer 205- `active_days` integer
README.md +5 −1
@@ -165,7 +165,7 @@ Example response:
165} 165}
166``` 166```
167 167
168Scheduled polling only runs when `ENABLE_SCHEDULER=true` and uses `DEFAULT_ACTOR`. 168Scheduled 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.
169 169
170For `git.sr.ht`, tracked repositories are configured via `GIT_TRACKED_REPOSITORIES`. Entries may be either: 170For `git.sr.ht`, tracked repositories are configured via `GIT_TRACKED_REPOSITORIES`. Entries may be either:
171 171
@@ -205,6 +205,9 @@ Example response:
205 "actor": "~your-user", 205 "actor": "~your-user",
206 "from": "2026-01-01", 206 "from": "2026-01-01",
207 "to": "2026-03-30", 207 "to": "2026-03-30",
208 "is_indexed": true,
209 "last_polled_at": "2026-04-11T18:05:00Z",
210 "indexing_state": "indexed",
208 "days": [ 211 "days": [
209 {"date": "2026-03-28", "count": 3, "score": 3.5}, 212 {"date": "2026-03-28", "count": 3, "score": 3.5},
210 {"date": "2026-03-29", "count": 0, "score": 0.0}, 213 {"date": "2026-03-29", "count": 0, "score": 0.0},
@@ -313,6 +316,7 @@ The SourceHut-specific assumptions are isolated to the service modules:
313 316
314- `git.sr.ht` polling is limited to repositories listed in `GIT_TRACKED_REPOSITORIES` 317- `git.sr.ht` polling is limited to repositories listed in `GIT_TRACKED_REPOSITORIES`
315- scheduled polling runs in-process, so it is not a distributed scheduler 318- scheduled polling runs in-process, so it is not a distributed scheduler
319- newly requested actors are indexed asynchronously, so the first public read may be empty until a scheduler or manual poll runs
316- alias management is config-driven; there is no alias CRUD API yet 320- alias management is config-driven; there is no alias CRUD API yet
317- current deployment model is trusted-operator V1, not a public multi-tenant service 321- current deployment model is trusted-operator V1, not a public multi-tenant service
318 322
alembic/versions/20260411_0002_tracked_actors.py added +39
@@ -0,0 +1,39 @@
1"""tracked actors for lazy indexing"""
2
3from __future__ import annotations
4
5from alembic import op
6import sqlalchemy as sa
7from sqlalchemy import inspect
8
9
10revision = "20260411_0002"
11down_revision = "20260409_0001"
12branch_labels = None
13depends_on = None
14
15
16def _table_names() -> set[str]:
17 return set(inspect(op.get_bind()).get_table_names())
18
19
20def upgrade() -> None:
21 if "tracked_actors" in _table_names():
22 return
23
24 op.create_table(
25 "tracked_actors",
26 sa.Column("id", sa.Integer(), primary_key=True),
27 sa.Column("actor", sa.String(length=255), nullable=False),
28 sa.Column("is_active", sa.Boolean(), nullable=False, server_default=sa.true()),
29 sa.Column("last_requested_at", sa.DateTime(timezone=True), nullable=True),
30 sa.Column("last_polled_at", sa.DateTime(timezone=True), nullable=True),
31 sa.Column("last_poll_status", sa.String(length=32), nullable=True),
32 sa.Column("last_poll_error", sa.Text(), nullable=True),
33 sa.UniqueConstraint("actor", name="uq_tracked_actor_actor"),
34 )
35
36
37def downgrade() -> None:
38 if "tracked_actors" in _table_names():
39 op.drop_table("tracked_actors")
src/srht_contrib/api/routes_contributions.py +10 −2
@@ -41,12 +41,16 @@ def get_contributions(
41 year: int | None = Query(default=None, ge=1970, le=3000), 41 year: int | None = Query(default=None, ge=1970, le=3000),
42 from_date: str | None = Query(default=None, alias="from"), 42 from_date: str | None = Query(default=None, alias="from"),
43 to_date: str | None = Query(default=None, alias="to"), 43 to_date: str | None = Query(default=None, alias="to"),
44 poller: PollerService = Depends(get_poller),
44 db: Session = Depends(get_db), 45 db: Session = Depends(get_db),
45 actor_identity_resolver: ActorIdentityResolver = Depends(get_actor_identity_resolver), 46 actor_identity_resolver: ActorIdentityResolver = Depends(get_actor_identity_resolver),
46) -> ContributionCalendarResponse: 47) -> ContributionCalendarResponse:
47 start, end = _resolve_range(year, from_date, to_date) 48 start, end = _resolve_range(year, from_date, to_date)
48 canonical_actor = actor_identity_resolver.canonicalize(actor, db=db) 49 canonical_actor = actor_identity_resolver.canonicalize(actor, db=db)
49 return ContributionAggregator().build_calendar(db, canonical_actor, start, end) 50 poller.track_actor_request(db, canonical_actor)
51 response = ContributionAggregator().build_calendar(db, canonical_actor, start, end)
52 db.commit()
53 return response
50 54
51 55
52@router.get("/{actor}/stats", response_model=ContributionStatsResponse) 56@router.get("/{actor}/stats", response_model=ContributionStatsResponse)
@@ -55,12 +59,16 @@ def get_contribution_stats(
55 year: int | None = Query(default=None, ge=1970, le=3000), 59 year: int | None = Query(default=None, ge=1970, le=3000),
56 from_date: str | None = Query(default=None, alias="from"), 60 from_date: str | None = Query(default=None, alias="from"),
57 to_date: str | None = Query(default=None, alias="to"), 61 to_date: str | None = Query(default=None, alias="to"),
62 poller: PollerService = Depends(get_poller),
58 db: Session = Depends(get_db), 63 db: Session = Depends(get_db),
59 actor_identity_resolver: ActorIdentityResolver = Depends(get_actor_identity_resolver), 64 actor_identity_resolver: ActorIdentityResolver = Depends(get_actor_identity_resolver),
60) -> ContributionStatsResponse: 65) -> ContributionStatsResponse:
61 start, end = _resolve_range(year, from_date, to_date) 66 start, end = _resolve_range(year, from_date, to_date)
62 canonical_actor = actor_identity_resolver.canonicalize(actor, db=db) 67 canonical_actor = actor_identity_resolver.canonicalize(actor, db=db)
63 return ContributionAggregator().build_stats(db, canonical_actor, start, end) 68 poller.track_actor_request(db, canonical_actor)
69 response = ContributionAggregator().build_stats(db, canonical_actor, start, end)
70 db.commit()
71 return response
64 72
65 73
66@router.post("/poll", response_model=PollResponse, dependencies=[Depends(require_api_key)]) 74@router.post("/poll", response_model=PollResponse, dependencies=[Depends(require_api_key)])
src/srht_contrib/jobs/poller.py +57 −2
@@ -7,9 +7,10 @@ from sqlalchemy import select
7from sqlalchemy.exc import IntegrityError 7from sqlalchemy.exc import IntegrityError
8from sqlalchemy.orm import Session 8from sqlalchemy.orm import Session
9 9
10from srht_contrib.models import ContributionEvent, SyncState, TrackedRepository 10from srht_contrib.models import ContributionEvent, SyncState, TrackedActor, TrackedRepository
11from srht_contrib.schemas import NormalizedEvent 11from srht_contrib.schemas import NormalizedEvent
12from srht_contrib.services.git import GitIngestionService 12from srht_contrib.services.git import GitIngestionService
13from srht_contrib.services.srht_client import SourceHutClientError
13from srht_contrib.services.todo import TodoIngestionService 14from srht_contrib.services.todo import TodoIngestionService
14from srht_contrib.utils.repositories import canonicalize_repository_name 15from srht_contrib.utils.repositories import canonicalize_repository_name
15 16
@@ -24,6 +25,52 @@ class PollerService:
24 self.git_service = git_service 25 self.git_service = git_service
25 26
26 def poll_all(self, db: Session, actor: str) -> int: 27 def poll_all(self, db: Session, actor: str) -> int:
28 self.track_actor_request(db, actor, update_last_requested=False)
29 try:
30 inserted = self._poll_actor(db, actor)
31 except Exception as exc:
32 db.rollback()
33 self._update_tracked_actor_poll_state(db, actor, status="error", error=str(exc))
34 db.commit()
35 raise
36
37 self._update_tracked_actor_poll_state(db, actor, status="indexed", error=None)
38 db.commit()
39 return inserted
40
41 def poll_tracked_actors(self, db: Session, default_actor: str | None = None) -> dict[str, int]:
42 if default_actor:
43 self.track_actor_request(db, default_actor, update_last_requested=False)
44 db.commit()
45
46 results: dict[str, int] = {}
47 actors = db.scalars(
48 select(TrackedActor.actor)
49 .where(TrackedActor.is_active.is_(True))
50 .order_by(TrackedActor.last_requested_at.is_(None), TrackedActor.last_requested_at.desc(), TrackedActor.actor)
51 ).all()
52 for actor in actors:
53 try:
54 results[actor] = self.poll_all(db, actor)
55 except SourceHutClientError:
56 logger.exception("Scheduled poll failed for actor=%s", actor)
57 except Exception:
58 logger.exception("Unexpected scheduled poll failure for actor=%s", actor)
59 return results
60
61 def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor:
62 tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor))
63 if tracked_actor is None:
64 tracked_actor = TrackedActor(actor=actor, is_active=True)
65 db.add(tracked_actor)
66
67 tracked_actor.is_active = True
68 if update_last_requested:
69 tracked_actor.last_requested_at = datetime.now(tz=UTC)
70 db.flush()
71 return tracked_actor
72
73 def _poll_actor(self, db: Session, actor: str) -> int:
27 inserted = 0 74 inserted = 0
28 inserted += self._poll_service(db, actor, self.todo_service.service_name, self.todo_service.fetch_recent_events) 75 inserted += self._poll_service(db, actor, self.todo_service.service_name, self.todo_service.fetch_recent_events)
29 self._sync_tracked_repositories(db, actor) 76 self._sync_tracked_repositories(db, actor)
@@ -38,7 +85,6 @@ class PollerService:
38 repositories=git_repositories, 85 repositories=git_repositories,
39 ), 86 ),
40 ) 87 )
41 db.commit()
42 return inserted 88 return inserted
43 89
44 def _poll_service(self, db: Session, actor: str, service_name: str, fetcher) -> int: 90 def _poll_service(self, db: Session, actor: str, service_name: str, fetcher) -> int:
@@ -117,3 +163,12 @@ class PollerService:
117 .order_by(TrackedRepository.repo_name) 163 .order_by(TrackedRepository.repo_name)
118 ).all() 164 ).all()
119 return list(rows) 165 return list(rows)
166
167 def _update_tracked_actor_poll_state(self, db: Session, actor: str, status: str, error: str | None) -> None:
168 tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
169 tracked_actor.last_poll_status = status
170 tracked_actor.last_poll_error = error
171 if status == "indexed":
172 tracked_actor.last_polled_at = datetime.now(tz=UTC)
173 db.add(tracked_actor)
174 db.flush()
src/srht_contrib/main.py +1 −1
@@ -91,7 +91,7 @@ def _scheduled_poll(app: FastAPI) -> None:
91 session_factory: sessionmaker[Session] = app.state.session_factory 91 session_factory: sessionmaker[Session] = app.state.session_factory
92 db = session_factory() 92 db = session_factory()
93 try: 93 try:
94 poller.poll_all(db, settings.default_actor) 94 poller.poll_tracked_actors(db, settings.default_actor)
95 finally: 95 finally:
96 db.close() 96 db.close()
97 97
src/srht_contrib/models.py +14 −1
@@ -2,7 +2,7 @@ from __future__ import annotations
2 2
3from datetime import datetime 3from datetime import datetime
4 4
5from sqlalchemy import JSON, DateTime, Float, Index, Integer, String, Text, UniqueConstraint 5from sqlalchemy import JSON, Boolean, DateTime, Float, Index, Integer, String, Text, UniqueConstraint
6from sqlalchemy.orm import Mapped, mapped_column 6from sqlalchemy.orm import Mapped, mapped_column
7 7
8from srht_contrib.db import Base 8from srht_contrib.db import Base
@@ -58,3 +58,16 @@ class ActorAlias(Base):
58 id: Mapped[int] = mapped_column(Integer, primary_key=True) 58 id: Mapped[int] = mapped_column(Integer, primary_key=True)
59 canonical_actor: Mapped[str] = mapped_column(String(255), nullable=False) 59 canonical_actor: Mapped[str] = mapped_column(String(255), nullable=False)
60 alias: Mapped[str] = mapped_column(String(255), nullable=False) 60 alias: Mapped[str] = mapped_column(String(255), nullable=False)
61
62
63class TrackedActor(Base):
64 __tablename__ = "tracked_actors"
65 __table_args__ = (UniqueConstraint("actor", name="uq_tracked_actor_actor"),)
66
67 id: Mapped[int] = mapped_column(Integer, primary_key=True)
68 actor: Mapped[str] = mapped_column(String(255), nullable=False)
69 is_active: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
70 last_requested_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
71 last_polled_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
72 last_poll_status: Mapped[str | None] = mapped_column(String(32), nullable=True)
73 last_poll_error: Mapped[str | None] = mapped_column(Text, nullable=True)
src/srht_contrib/schemas.py +9 −2
@@ -1,6 +1,7 @@
1from __future__ import annotations 1from __future__ import annotations
2 2
3from datetime import date, datetime 3from datetime import date, datetime
4from typing import Literal
4 5
5from pydantic import BaseModel, Field, model_validator 6from pydantic import BaseModel, Field, model_validator
6 7
@@ -23,7 +24,13 @@ class ContributionDay(BaseModel):
23 score: float 24 score: float
24 25
25 26
26class ContributionCalendarResponse(BaseModel): 27class ContributionIndexMetadata(BaseModel):
28 is_indexed: bool
29 last_polled_at: datetime | None = None
30 indexing_state: Literal["pending", "indexed", "error"]
31
32
33class ContributionCalendarResponse(ContributionIndexMetadata):
27 actor: str 34 actor: str
28 from_date: date = Field(alias="from") 35 from_date: date = Field(alias="from")
29 to_date: date = Field(alias="to") 36 to_date: date = Field(alias="to")
@@ -32,7 +39,7 @@ class ContributionCalendarResponse(BaseModel):
32 model_config = {"populate_by_name": True} 39 model_config = {"populate_by_name": True}
33 40
34 41
35class ContributionStatsResponse(BaseModel): 42class ContributionStatsResponse(ContributionIndexMetadata):
36 actor: str 43 actor: str
37 from_date: date = Field(alias="from") 44 from_date: date = Field(alias="from")
38 to_date: date = Field(alias="to") 45 to_date: date = Field(alias="to")
src/srht_contrib/services/aggregator.py +37 −3
@@ -6,8 +6,13 @@ from datetime import date
6from sqlalchemy import func, select 6from sqlalchemy import func, select
7from sqlalchemy.orm import Session 7from sqlalchemy.orm import Session
8 8
9from srht_contrib.models import ContributionEvent 9from srht_contrib.models import ContributionEvent, TrackedActor
10from srht_contrib.schemas import ContributionCalendarResponse, ContributionDay, ContributionStatsResponse 10from srht_contrib.schemas import (
11 ContributionCalendarResponse,
12 ContributionDay,
13 ContributionIndexMetadata,
14 ContributionStatsResponse,
15)
11from srht_contrib.utils.dates import date_range, date_to_utc_bounds 16from srht_contrib.utils.dates import date_range, date_to_utc_bounds
12 17
13 18
@@ -30,7 +35,14 @@ class ContributionAggregator:
30 ) 35 )
31 for day in date_range(start, end) 36 for day in date_range(start, end)
32 ] 37 ]
33 return ContributionCalendarResponse(actor=actor, from_date=start, to_date=end, days=days) 38 metadata = self._index_metadata(db, actor)
39 return ContributionCalendarResponse(
40 actor=actor,
41 from_date=start,
42 to_date=end,
43 days=days,
44 **metadata.model_dump(),
45 )
34 46
35 def build_stats(self, db: Session, actor: str, start: date, end: date) -> ContributionStatsResponse: 47 def build_stats(self, db: Session, actor: str, start: date, end: date) -> ContributionStatsResponse:
36 calendar = self.build_calendar(db, actor, start, end) 48 calendar = self.build_calendar(db, actor, start, end)
@@ -47,6 +59,28 @@ class ContributionAggregator:
47 active_days=len(active_days), 59 active_days=len(active_days),
48 longest_streak=max(streaks, default=0), 60 longest_streak=max(streaks, default=0),
49 current_streak=current_streak, 61 current_streak=current_streak,
62 is_indexed=calendar.is_indexed,
63 last_polled_at=calendar.last_polled_at,
64 indexing_state=calendar.indexing_state,
65 )
66
67 def _index_metadata(self, db: Session, actor: str) -> ContributionIndexMetadata:
68 tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor))
69 has_indexed_events = db.scalar(select(ContributionEvent.id).where(ContributionEvent.actor == actor).limit(1)) is not None
70 last_poll_status = tracked_actor.last_poll_status if tracked_actor is not None else None
71 is_indexed = has_indexed_events or (tracked_actor is not None and tracked_actor.last_polled_at is not None)
72
73 if last_poll_status == "error":
74 indexing_state = "error"
75 elif is_indexed:
76 indexing_state = "indexed"
77 else:
78 indexing_state = "pending"
79
80 return ContributionIndexMetadata(
81 is_indexed=is_indexed,
82 last_polled_at=tracked_actor.last_polled_at if tracked_actor is not None else None,
83 indexing_state=indexing_state,
50 ) 84 )
51 85
52 def _query_daily_aggregates(self, db: Session, actor: str, start: date, end: date) -> list[DailyAggregate]: 86 def _query_daily_aggregates(self, db: Session, actor: str, start: date, end: date) -> list[DailyAggregate]:
tests/test_contributions_api.py +20 −1
@@ -1,9 +1,10 @@
1from datetime import UTC, datetime 1from datetime import UTC, datetime
2 2
3from fastapi.testclient import TestClient 3from fastapi.testclient import TestClient
4from sqlalchemy import select
4 5
5from srht_contrib.main import create_app 6from srht_contrib.main import create_app
6from srht_contrib.models import ContributionEvent 7from srht_contrib.models import ContributionEvent, TrackedActor
7 8
8 9
9def test_read_only_contribution_routes_are_public_and_write_routes_require_api_key(settings, db_engine, session_factory) -> None: 10def 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
45 response = client.get("/api/contributions/~ccleberg?from=2026-03-28&to=2026-03-30") 46 response = client.get("/api/contributions/~ccleberg?from=2026-03-28&to=2026-03-30")
46 47
47 assert response.status_code == 200 48 assert response.status_code == 200
49 assert response.json()["is_indexed"] is True
50 assert response.json()["indexing_state"] == "indexed"
48 assert response.json()["days"] == [ 51 assert response.json()["days"] == [
49 {"date": "2026-03-28", "count": 0, "score": 0.0}, 52 {"date": "2026-03-28", "count": 0, "score": 0.0},
50 {"date": "2026-03-29", "count": 0, "score": 0.0}, 53 {"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
52 ] 55 ]
53 56
54 57
58def test_public_read_registers_actor_for_lazy_indexing(client: TestClient, db_session) -> None:
59 response = client.get("/api/contributions/~ccleberg?from=2026-03-28&to=2026-03-30")
60
61 tracked_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg"))
62
63 assert response.status_code == 200
64 assert response.json()["is_indexed"] is False
65 assert response.json()["indexing_state"] == "pending"
66 assert response.json()["last_polled_at"] is None
67 assert tracked_actor is not None
68 assert tracked_actor.is_active is True
69 assert tracked_actor.last_requested_at is not None
70
71
55def test_contribution_stats_api(client: TestClient, db_session) -> None: 72def test_contribution_stats_api(client: TestClient, db_session) -> None:
56 db_session.add_all( 73 db_session.add_all(
57 [ 74 [
@@ -88,6 +105,8 @@ def test_contribution_stats_api(client: TestClient, db_session) -> None:
88 assert response.json()["total_score"] == 1.5 105 assert response.json()["total_score"] == 1.5
89 assert response.json()["longest_streak"] == 2 106 assert response.json()["longest_streak"] == 2
90 assert response.json()["current_streak"] == 2 107 assert response.json()["current_streak"] == 2
108 assert response.json()["is_indexed"] is True
109 assert response.json()["indexing_state"] == "indexed"
91 110
92 111
93def test_invalid_date_input_returns_400(client: TestClient) -> None: 112def test_invalid_date_input_returns_400(client: TestClient) -> None:
tests/test_ingestion.py +29 −1
@@ -4,7 +4,7 @@ from sqlalchemy import select
4 4
5from srht_contrib.config import Settings 5from srht_contrib.config import Settings
6from srht_contrib.jobs.poller import PollerService 6from srht_contrib.jobs.poller import PollerService
7from srht_contrib.models import SyncState, TrackedRepository 7from srht_contrib.models import SyncState, TrackedActor, TrackedRepository
8from srht_contrib.schemas import NormalizedEvent 8from srht_contrib.schemas import NormalizedEvent
9from srht_contrib.services.git import GitIngestionService, GitPollResult 9from srht_contrib.services.git import GitIngestionService, GitPollResult
10from srht_contrib.services.todo import TodoIngestionService, TodoPollResult 10from srht_contrib.services.todo import TodoIngestionService, TodoPollResult
@@ -335,3 +335,31 @@ def test_sync_overlap_reuses_cursor_window_and_suppresses_duplicates(db_session)
335 assert state is not None 335 assert state is not None
336 assert len(todo_service.calls) == 2 336 assert len(todo_service.calls) == 2
337 assert todo_service.calls[1].isoformat() == "2026-03-30T00:00:00+00:00" 337 assert todo_service.calls[1].isoformat() == "2026-03-30T00:00:00+00:00"
338
339
340def test_scheduled_poll_polls_known_actors_and_seeds_default_actor(db_session) -> None:
341 event = NormalizedEvent(
342 service="todo",
343 event_type="ticket_created",
344 actor="~known",
345 repo_name="todo",
346 resource_id="123",
347 external_uid="todo:event:known:created:123",
348 occurred_at=datetime(2026, 3, 30, 10, 0, tzinfo=UTC),
349 weight=1.0,
350 raw_payload_json=None,
351 )
352 todo_service = RecordingTodoService(events_by_call=[[], [event]])
353 poller = PollerService(todo_service=todo_service, git_service=EmptyGitService())
354
355 db_session.add(TrackedActor(actor="~known", is_active=True))
356 db_session.commit()
357
358 results = poller.poll_tracked_actors(db_session, default_actor="~default")
359
360 tracked_actors = db_session.scalars(select(TrackedActor).order_by(TrackedActor.actor)).all()
361
362 assert results == {"~default": 0, "~known": 1}
363 assert [actor.actor for actor in tracked_actors] == ["~default", "~known"]
364 assert all(actor.last_poll_status == "indexed" for actor in tracked_actors)
365 assert all(actor.last_polled_at is not None for actor in tracked_actors)
tests/test_migrations.py +3
@@ -22,6 +22,7 @@ def test_alembic_upgrade_creates_schema(tmp_path) -> None:
22 "alembic_version", 22 "alembic_version",
23 "contribution_events", 23 "contribution_events",
24 "sync_states", 24 "sync_states",
25 "tracked_actors",
25 "tracked_repositories", 26 "tracked_repositories",
26 ] 27 ]
27 28
@@ -115,6 +116,7 @@ def test_alembic_upgrade_adopts_legacy_schema(tmp_path) -> None:
115 assert columns["actor"]["nullable"] is False 116 assert columns["actor"]["nullable"] is False
116 assert "uq_tracked_repository_service_actor_name" in unique_constraints 117 assert "uq_tracked_repository_service_actor_name" in unique_constraints
117 assert actor == Settings().default_actor 118 assert actor == Settings().default_actor
119 assert "tracked_actors" in inspector.get_table_names()
118 120
119 121
120def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) -> None: 122def 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) ->
128 130
129 inspector = inspect(create_engine(database_url)) 131 inspector = inspect(create_engine(database_url))
130 assert "actor_aliases" in inspector.get_table_names() 132 assert "actor_aliases" in inspector.get_table_names()
133 assert "tracked_actors" in inspector.get_table_names()
tests/test_polling_api.py +21 −1
@@ -1,10 +1,11 @@
1from datetime import UTC, datetime 1from datetime import UTC, datetime
2 2
3from fastapi.testclient import TestClient 3from fastapi.testclient import TestClient
4from sqlalchemy import select
4 5
5from srht_contrib.config import Settings 6from srht_contrib.config import Settings
6from srht_contrib.main import create_app 7from srht_contrib.main import create_app
7from srht_contrib.models import ContributionEvent 8from srht_contrib.models import ContributionEvent, TrackedActor
8from srht_contrib.services.srht_client import SourceHutClientError 9from srht_contrib.services.srht_client import SourceHutClientError
9 10
10 11
@@ -19,6 +20,16 @@ class InsertingPoller:
19 self.todo_service = service 20 self.todo_service = service
20 self.git_service = service 21 self.git_service = service
21 22
23 def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True):
24 tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor))
25 if tracked_actor is None:
26 tracked_actor = TrackedActor(actor=actor, is_active=True)
27 db.add(tracked_actor)
28 if update_last_requested:
29 tracked_actor.last_requested_at = datetime(2026, 3, 30, 9, 0, tzinfo=UTC)
30 db.flush()
31 return tracked_actor
32
22 def poll_all(self, db, actor: str) -> int: 33 def poll_all(self, db, actor: str) -> int:
23 db.add( 34 db.add(
24 ContributionEvent( 35 ContributionEvent(
@@ -43,6 +54,14 @@ class FailingPoller:
43 self.todo_service = service 54 self.todo_service = service
44 self.git_service = service 55 self.git_service = service
45 56
57 def track_actor_request(self, db, actor: str, *, update_last_requested: bool = True):
58 tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor))
59 if tracked_actor is None:
60 tracked_actor = TrackedActor(actor=actor, is_active=True)
61 db.add(tracked_actor)
62 db.flush()
63 return tracked_actor
64
46 def poll_all(self, db, actor: str) -> int: 65 def poll_all(self, db, actor: str) -> int:
47 raise SourceHutClientError("boom") 66 raise SourceHutClientError("boom")
48 67
@@ -58,6 +77,7 @@ def test_manual_poll_uses_same_database_session(settings: Settings, db_engine, s
58 assert poll_response.status_code == 200 77 assert poll_response.status_code == 200
59 assert poll_response.json()["inserted_events"] == 1 78 assert poll_response.json()["inserted_events"] == 1
60 assert calendar_response.status_code == 200 79 assert calendar_response.status_code == 200
80 assert calendar_response.json()["is_indexed"] is True
61 assert calendar_response.json()["days"] == [{"date": "2026-03-30", "count": 1, "score": 1.0}] 81 assert calendar_response.json()["days"] == [{"date": "2026-03-30", "count": 1, "score": 1.0}]
62 82
63 83