krz/hutch-stats

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

Commit e1a3233ec2

e1a3233ec2204055faa8dadb18cb2efa12cf0e17

parent: aebe493dcd

Verified · cmc

cmc <hello@cleberg.net> · 2026-04-11 22:39 UTC

feat: add resumable historical backfill for actor activity

Layout: unified · split

API.md +24 −1
@@ -44,6 +44,7 @@ Contribution 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- Public contribution reads also register the actor for background indexing. A first lookup may therefore return an empty graph while the scheduler catches up.
47- Incremental indexing and historical backfill are separate. An actor can be recently indexed without being fully backfilled yet.
47 48
48Background polling: 49Background polling:
49 50
@@ -110,6 +111,7 @@ Behavior notes:
110- This endpoint resolves aliases to a canonical actor before querying data. 111- This endpoint resolves aliases to a canonical actor before querying data.
111- This endpoint also registers the actor for background indexing and updates the actor's `last_requested_at` timestamp. 112- This endpoint also registers the actor for background indexing and updates the actor's `last_requested_at` timestamp.
112- The response is always immediate; it does not wait for SourceHut polling to finish. 113- The response is always immediate; it does not wait for SourceHut polling to finish.
114- Historical backfill runs in bounded background batches and may take multiple scheduler passes to complete.
113 115
114Example by year: 116Example by year:
115 117
@@ -133,6 +135,9 @@ Response `200 OK`:
133 "is_indexed": false, 135 "is_indexed": false,
134 "last_polled_at": null, 136 "last_polled_at": null,
135 "indexing_state": "pending", 137 "indexing_state": "pending",
138 "is_backfilled": false,
139 "backfill_state": "in_progress",
140 "backfill_completed_at": null,
136 "days": [ 141 "days": [
137 { "date": "2026-03-01", "count": 0, "score": 0.0 }, 142 { "date": "2026-03-01", "count": 0, "score": 0.0 },
138 { "date": "2026-03-02", "count": 3, "score": 2.5 } 143 { "date": "2026-03-02", "count": 3, "score": 2.5 }
@@ -145,9 +150,12 @@ Response fields:
145- `actor` string: canonical actor after alias resolution 150- `actor` string: canonical actor after alias resolution
146- `from` string: inclusive start date 151- `from` string: inclusive start date
147- `to` string: inclusive end date 152- `to` string: inclusive end date
148- `is_indexed` boolean: whether the service has already indexed activity for this actor 153- `is_indexed` boolean: whether the service has already completed at least one successful recent/incremental poll for this actor
149- `last_polled_at` string or `null`: most recent successful poll time, if any 154- `last_polled_at` string or `null`: most recent successful poll time, if any
150- `indexing_state` string: one of `pending`, `indexed`, or `error` 155- `indexing_state` string: one of `pending`, `indexed`, or `error`
156- `is_backfilled` boolean: whether historical backfill has completed for this actor
157- `backfill_state` string: one of `pending`, `in_progress`, `completed`, or `error`
158- `backfill_completed_at` string or `null`: when full historical backfill completed, if it has
151- `days` array: 159- `days` array:
152 - `date` string `YYYY-MM-DD` 160 - `date` string `YYYY-MM-DD`
153 - `count` integer contribution count for the day 161 - `count` integer contribution count for the day
@@ -159,6 +167,13 @@ Indexing state semantics:
159- `indexed`: at least one successful poll has completed for the actor 167- `indexed`: at least one successful poll has completed for the actor
160- `error`: the most recent poll attempt for the actor failed 168- `error`: the most recent poll attempt for the actor failed
161 169
170Backfill state semantics:
171
172- `pending`: the actor has not started historical backfill yet
173- `in_progress`: historical backfill is actively progressing in bounded background batches
174- `completed`: historical backfill has completed for all supported services
175- `error`: the most recent backfill attempt failed
176
162Possible errors: 177Possible errors:
163 178
164- `400 Bad Request` for invalid or conflicting date input 179- `400 Bad Request` for invalid or conflicting date input
@@ -193,6 +208,7 @@ Behavior notes:
193 208
194- This endpoint has the same actor-registration and alias-resolution behavior as the calendar endpoint. 209- This endpoint has the same actor-registration and alias-resolution behavior as the calendar endpoint.
195- This endpoint returns immediately and does not block on SourceHut polling. 210- This endpoint returns immediately and does not block on SourceHut polling.
211- This endpoint also reflects historical backfill state so clients can distinguish recent indexing from complete history.
196 212
197Example: 213Example:
198 214
@@ -210,6 +226,9 @@ Response `200 OK`:
210 "is_indexed": true, 226 "is_indexed": true,
211 "last_polled_at": "2026-04-11T18:05:00Z", 227 "last_polled_at": "2026-04-11T18:05:00Z",
212 "indexing_state": "indexed", 228 "indexing_state": "indexed",
229 "is_backfilled": false,
230 "backfill_state": "in_progress",
231 "backfill_completed_at": null,
213 "total_events": 126, 232 "total_events": 126,
214 "total_score": 116.75, 233 "total_score": 116.75,
215 "active_days": 14, 234 "active_days": 14,
@@ -226,6 +245,9 @@ Response fields:
226- `is_indexed` boolean 245- `is_indexed` boolean
227- `last_polled_at` string or `null` 246- `last_polled_at` string or `null`
228- `indexing_state` string 247- `indexing_state` string
248- `is_backfilled` boolean
249- `backfill_state` string
250- `backfill_completed_at` string or `null`
229- `total_events` integer 251- `total_events` integer
230- `total_score` float 252- `total_score` float
231- `active_days` integer 253- `active_days` integer
@@ -275,6 +297,7 @@ Response fields:
275Behavior notes: 297Behavior notes:
276 298
277- Manual polling also updates the actor's indexing metadata. 299- Manual polling also updates the actor's indexing metadata.
300- Manual polling also advances historical backfill by one bounded batch per supported service.
278- Git polling auto-discovers the actor's owned repositories and unions in any configured tracked repositories. 301- Git polling auto-discovers the actor's owned repositories and unions in any configured tracked repositories.
279 302
280Possible errors: 303Possible errors:
README.md +13 −3
@@ -14,7 +14,7 @@ The current V1 is intentionally narrow and production-oriented:
14 14
15## What It Does 15## What It Does
16 16
17The service collects SourceHut activity from one or more sr.ht GraphQL services, turns those records into a canonical event shape, aggregates activity by day, and returns zero-filled calendar ranges so the client never has to patch missing dates. 17The service collects SourceHut activity from one or more sr.ht GraphQL services, turns those records into a canonical event shape, aggregates activity by day, and returns zero-filled calendar ranges so the client never has to patch missing dates. It performs both recent incremental polling and bounded historical backfill.
18 18
19Example use cases: 19Example use cases:
20 20
@@ -28,7 +28,7 @@ The code is split into small, testable layers:
28 28
29- `src/srht_contrib/config.py`: environment-driven settings and event weights 29- `src/srht_contrib/config.py`: environment-driven settings and event weights
30- `src/srht_contrib/db.py`: SQLAlchemy engine/session setup and app-scoped DB access 30- `src/srht_contrib/db.py`: SQLAlchemy engine/session setup and app-scoped DB access
31- `src/srht_contrib/models.py`: ORM models for normalized events, sync state, aliases, and tracked repos 31- `src/srht_contrib/models.py`: ORM models for normalized events, sync state, aliases, tracked actors, and backfill state
32- `src/srht_contrib/services/srht_client.py`: generic SourceHut GraphQL client with error handling and simple retries 32- `src/srht_contrib/services/srht_client.py`: generic SourceHut GraphQL client with error handling and simple retries
33- `src/srht_contrib/services/todo.py`: `todo.sr.ht` ingestion and normalization 33- `src/srht_contrib/services/todo.py`: `todo.sr.ht` ingestion and normalization
34- `src/srht_contrib/services/git.py`: `git.sr.ht` tracked-repository commit ingestion 34- `src/srht_contrib/services/git.py`: `git.sr.ht` tracked-repository commit ingestion
@@ -165,7 +165,7 @@ Example response:
165} 165}
166``` 166```
167 167
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. 168Scheduled polling only runs when `ENABLE_SCHEDULER=true`. The scheduler seeds `DEFAULT_ACTOR` as an initial known actor, runs one poll immediately at startup, and public contribution reads register additional actors for later background polling and historical backfill.
169 169
170For `git.sr.ht`, owned repositories are auto-discovered for the actor. `GIT_TRACKED_REPOSITORIES` can still be used to union in extra repositories. Entries may be either: 170For `git.sr.ht`, owned repositories are auto-discovered for the actor. `GIT_TRACKED_REPOSITORIES` can still be used to union in extra repositories. Entries may be either:
171 171
@@ -208,6 +208,9 @@ Example response:
208 "is_indexed": true, 208 "is_indexed": true,
209 "last_polled_at": "2026-04-11T18:05:00Z", 209 "last_polled_at": "2026-04-11T18:05:00Z",
210 "indexing_state": "indexed", 210 "indexing_state": "indexed",
211 "is_backfilled": false,
212 "backfill_state": "in_progress",
213 "backfill_completed_at": null,
211 "days": [ 214 "days": [
212 {"date": "2026-03-28", "count": 3, "score": 3.5}, 215 {"date": "2026-03-28", "count": 3, "score": 3.5},
213 {"date": "2026-03-29", "count": 0, "score": 0.0}, 216 {"date": "2026-03-29", "count": 0, "score": 0.0},
@@ -229,6 +232,12 @@ Example response:
229 "actor": "~your-user", 232 "actor": "~your-user",
230 "from": "2026-01-01", 233 "from": "2026-01-01",
231 "to": "2026-12-31", 234 "to": "2026-12-31",
235 "is_indexed": true,
236 "last_polled_at": "2026-04-11T18:05:00Z",
237 "indexing_state": "indexed",
238 "is_backfilled": false,
239 "backfill_state": "in_progress",
240 "backfill_completed_at": null,
232 "total_events": 42, 241 "total_events": 42,
233 "total_score": 37.5, 242 "total_score": 37.5,
234 "active_days": 18, 243 "active_days": 18,
@@ -317,6 +326,7 @@ The SourceHut-specific assumptions are isolated to the service modules:
317- `git.sr.ht` polling assumes the actor's repositories are discoverable through the SourceHut GraphQL API 326- `git.sr.ht` polling assumes the actor's repositories are discoverable through the SourceHut GraphQL API
318- scheduled polling runs in-process, so it is not a distributed scheduler 327- 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 328- newly requested actors are indexed asynchronously, so the first public read may be empty until a scheduler or manual poll runs
329- full historical backfill can take many scheduler passes for active users because it runs in bounded batches
320- alias management is config-driven; there is no alias CRUD API yet 330- alias management is config-driven; there is no alias CRUD API yet
321- current deployment model is trusted-operator V1, not a public multi-tenant service 331- current deployment model is trusted-operator V1, not a public multi-tenant service
322 332
alembic/versions/20260411_0003_backfill_state.py added +55
@@ -0,0 +1,55 @@
1"""actor and service backfill state"""
2
3from __future__ import annotations
4
5from alembic import op
6import sqlalchemy as sa
7from sqlalchemy import inspect
8
9
10revision = "20260411_0003"
11down_revision = "20260411_0002"
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 _column_names(table_name: str) -> set[str]:
21 return {column["name"] for column in inspect(op.get_bind()).get_columns(table_name)}
22
23
24def upgrade() -> None:
25 if "tracked_actors" in _table_names():
26 columns = _column_names("tracked_actors")
27 if "backfill_status" not in columns:
28 op.add_column("tracked_actors", sa.Column("backfill_status", sa.String(length=32), nullable=True))
29 op.execute(sa.text("UPDATE tracked_actors SET backfill_status = 'pending' WHERE backfill_status IS NULL"))
30 if "backfill_started_at" not in columns:
31 op.add_column("tracked_actors", sa.Column("backfill_started_at", sa.DateTime(timezone=True), nullable=True))
32 if "backfill_completed_at" not in columns:
33 op.add_column("tracked_actors", sa.Column("backfill_completed_at", sa.DateTime(timezone=True), nullable=True))
34 if "last_backfill_error" not in columns:
35 op.add_column("tracked_actors", sa.Column("last_backfill_error", sa.Text(), nullable=True))
36
37 if "service_backfill_states" not in _table_names():
38 op.create_table(
39 "service_backfill_states",
40 sa.Column("id", sa.Integer(), primary_key=True),
41 sa.Column("actor", sa.String(length=255), nullable=False),
42 sa.Column("service", sa.String(length=32), nullable=False),
43 sa.Column("cursor_json", sa.JSON(), nullable=True),
44 sa.Column("status", sa.String(length=32), nullable=False),
45 sa.Column("started_at", sa.DateTime(timezone=True), nullable=True),
46 sa.Column("completed_at", sa.DateTime(timezone=True), nullable=True),
47 sa.Column("last_error", sa.Text(), nullable=True),
48 sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False),
49 sa.UniqueConstraint("actor", "service", name="uq_service_backfill_state_actor_service"),
50 )
51
52
53def downgrade() -> None:
54 if "service_backfill_states" in _table_names():
55 op.drop_table("service_backfill_states")
src/srht_contrib/jobs/poller.py +89 −2
@@ -7,7 +7,7 @@ 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, TrackedActor, TrackedRepository 10from srht_contrib.models import ContributionEvent, ServiceBackfillState, 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.srht_client import SourceHutClientError
@@ -35,6 +35,7 @@ class PollerService:
35 raise 35 raise
36 36
37 self._update_tracked_actor_poll_state(db, actor, status="indexed", error=None) 37 self._update_tracked_actor_poll_state(db, actor, status="indexed", error=None)
38 inserted += self._run_backfill_batches(db, actor)
38 db.commit() 39 db.commit()
39 return inserted 40 return inserted
40 41
@@ -61,7 +62,7 @@ class PollerService:
61 def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor: 62 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 tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor))
63 if tracked_actor is None: 64 if tracked_actor is None:
64 tracked_actor = TrackedActor(actor=actor, is_active=True) 65 tracked_actor = TrackedActor(actor=actor, is_active=True, backfill_status="pending")
65 db.add(tracked_actor) 66 db.add(tracked_actor)
66 67
67 tracked_actor.is_active = True 68 tracked_actor.is_active = True
@@ -172,3 +173,89 @@ class PollerService:
172 tracked_actor.last_polled_at = datetime.now(tz=UTC) 173 tracked_actor.last_polled_at = datetime.now(tz=UTC)
173 db.add(tracked_actor) 174 db.add(tracked_actor)
174 db.flush() 175 db.flush()
176
177 def _run_backfill_batches(self, db: Session, actor: str) -> int:
178 tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
179 if tracked_actor.backfill_status == "completed":
180 return 0
181
182 if tracked_actor.backfill_started_at is None:
183 tracked_actor.backfill_started_at = datetime.now(tz=UTC)
184 tracked_actor.backfill_status = "in_progress"
185 tracked_actor.last_backfill_error = None
186 db.add(tracked_actor)
187 db.flush()
188 total_inserted = 0
189
190 services = [
191 (self.todo_service.service_name, self.todo_service.fetch_backfill_batch),
192 (self.git_service.service_name, self.git_service.fetch_backfill_batch),
193 ]
194 all_complete = True
195 for service_name, fetcher in services:
196 state = db.scalar(
197 select(ServiceBackfillState)
198 .where(ServiceBackfillState.actor == actor)
199 .where(ServiceBackfillState.service == service_name)
200 )
201 if state is None:
202 state = ServiceBackfillState(
203 actor=actor,
204 service=service_name,
205 cursor_json=None,
206 status="pending",
207 started_at=None,
208 completed_at=None,
209 last_error=None,
210 updated_at=datetime.now(tz=UTC),
211 )
212 db.add(state)
213 db.flush()
214
215 if state.status == "completed":
216 continue
217
218 all_complete = False
219 if state.started_at is None:
220 state.started_at = datetime.now(tz=UTC)
221 state.status = "in_progress"
222 state.updated_at = datetime.now(tz=UTC)
223 try:
224 result = fetcher(actor=actor, cursor_state=state.cursor_json)
225 inserted = self._insert_events(db, result.events)
226 total_inserted += inserted
227 state.cursor_json = result.cursor_state
228 state.last_error = None
229 state.updated_at = datetime.now(tz=UTC)
230 if result.complete:
231 state.status = "completed"
232 state.completed_at = datetime.now(tz=UTC)
233 logger.info("Backfill complete for service=%s actor=%s inserted=%s", service_name, actor, inserted)
234 else:
235 logger.info("Backfill batch complete for service=%s actor=%s inserted=%s", service_name, actor, inserted)
236 except Exception as exc:
237 state.status = "error"
238 state.last_error = str(exc)
239 state.updated_at = datetime.now(tz=UTC)
240 tracked_actor.backfill_status = "error"
241 tracked_actor.last_backfill_error = str(exc)
242 db.add(state)
243 db.add(tracked_actor)
244 db.flush()
245 raise
246
247 db.add(state)
248 db.flush()
249
250 completed = db.scalars(
251 select(ServiceBackfillState.status).where(ServiceBackfillState.actor == actor)
252 ).all()
253 if completed and all(status == "completed" for status in completed):
254 tracked_actor.backfill_status = "completed"
255 tracked_actor.backfill_completed_at = datetime.now(tz=UTC)
256 tracked_actor.last_backfill_error = None
257 elif tracked_actor.backfill_status != "error":
258 tracked_actor.backfill_status = "in_progress"
259 db.add(tracked_actor)
260 db.flush()
261 return total_inserted
src/srht_contrib/models.py +19
@@ -71,3 +71,22 @@ class TrackedActor(Base):
71 last_polled_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) 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) 73 last_poll_error: Mapped[str | None] = mapped_column(Text, nullable=True)
74 backfill_status: Mapped[str] = mapped_column(String(32), nullable=False, default="pending")
75 backfill_started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
76 backfill_completed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
77 last_backfill_error: Mapped[str | None] = mapped_column(Text, nullable=True)
78
79
80class ServiceBackfillState(Base):
81 __tablename__ = "service_backfill_states"
82 __table_args__ = (UniqueConstraint("actor", "service", name="uq_service_backfill_state_actor_service"),)
83
84 id: Mapped[int] = mapped_column(Integer, primary_key=True)
85 actor: Mapped[str] = mapped_column(String(255), nullable=False)
86 service: Mapped[str] = mapped_column(String(32), nullable=False)
87 cursor_json: Mapped[dict | None] = mapped_column(JSON, nullable=True)
88 status: Mapped[str] = mapped_column(String(32), nullable=False, default="pending")
89 started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
90 completed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
91 last_error: Mapped[str | None] = mapped_column(Text, nullable=True)
92 updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False)
src/srht_contrib/schemas.py +3
@@ -28,6 +28,9 @@ class ContributionIndexMetadata(BaseModel):
28 is_indexed: bool 28 is_indexed: bool
29 last_polled_at: datetime | None = None 29 last_polled_at: datetime | None = None
30 indexing_state: Literal["pending", "indexed", "error"] 30 indexing_state: Literal["pending", "indexed", "error"]
31 is_backfilled: bool = False
32 backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending"
33 backfill_completed_at: datetime | None = None
31 34
32 35
33class ContributionCalendarResponse(ContributionIndexMetadata): 36class ContributionCalendarResponse(ContributionIndexMetadata):
src/srht_contrib/services/aggregator.py +3
@@ -81,6 +81,9 @@ class ContributionAggregator:
81 is_indexed=is_indexed, 81 is_indexed=is_indexed,
82 last_polled_at=tracked_actor.last_polled_at if tracked_actor is not None else None, 82 last_polled_at=tracked_actor.last_polled_at if tracked_actor is not None else None,
83 indexing_state=indexing_state, 83 indexing_state=indexing_state,
84 is_backfilled=(tracked_actor.backfill_status == "completed") if tracked_actor is not None else False,
85 backfill_state=(tracked_actor.backfill_status if tracked_actor is not None else "pending"),
86 backfill_completed_at=tracked_actor.backfill_completed_at if tracked_actor is not None else None,
84 ) 87 )
85 88
86 def _query_daily_aggregates(self, db: Session, actor: str, start: date, end: date) -> list[DailyAggregate]: 89 def _query_daily_aggregates(self, db: Session, actor: str, start: date, end: date) -> list[DailyAggregate]:
src/srht_contrib/services/git.py +90
@@ -8,6 +8,7 @@ from typing import Any
8from srht_contrib.config import Settings 8from srht_contrib.config import Settings
9from srht_contrib.schemas import NormalizedEvent 9from srht_contrib.schemas import NormalizedEvent
10from srht_contrib.services.srht_client import SourceHutGraphQLClient 10from srht_contrib.services.srht_client import SourceHutGraphQLClient
11from srht_contrib.services.types import BackfillBatchResult
11from srht_contrib.utils.dates import ensure_utc, parse_datetime 12from srht_contrib.utils.dates import ensure_utc, parse_datetime
12from srht_contrib.utils.identity import ActorIdentityResolver 13from srht_contrib.utils.identity import ActorIdentityResolver
13 14
@@ -117,6 +118,95 @@ class GitIngestionService:
117 ) 118 )
118 return repositories 119 return repositories
119 120
121 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
122 state = {
123 "discovery_cursor": None,
124 "discovery_complete": False,
125 "repository_queue": sorted(
126 {
127 self._canonical_repository_name(actor, repository)
128 for repository in self.settings.git_tracked_repositories
129 }
130 ),
131 "current_repository": None,
132 }
133 if cursor_state:
134 state.update(cursor_state)
135
136 if not state["discovery_complete"]:
137 data = self.client.execute(
138 USER_REPOSITORIES_QUERY,
139 {"username": actor.lstrip("~"), "cursor": state["discovery_cursor"]},
140 )
141 user = data.get("user") or {}
142 repositories_page = user.get("repositories") or {}
143 results = repositories_page.get("results") or []
144 state["discovery_cursor"] = repositories_page.get("cursor")
145 state["discovery_complete"] = not bool(state["discovery_cursor"])
146 known = set(state["repository_queue"])
147 current_repository = state.get("current_repository")
148 if current_repository:
149 known.add(current_repository["name"])
150 for repository in results:
151 if not isinstance(repository, dict):
152 continue
153 name = repository.get("name")
154 repository_owner = ((repository.get("owner") or {}).get("canonicalName") or actor).strip()
155 if not name or not repository_owner:
156 continue
157 canonical_name = f"{repository_owner}/{name}"
158 if canonical_name not in known:
159 state["repository_queue"].append(canonical_name)
160 known.add(canonical_name)
161 state["repository_queue"] = sorted(state["repository_queue"])
162 logger.info(
163 "git backfill discovery actor=%s page_count=%s queue=%s next_cursor=%s",
164 actor,
165 len(results),
166 len(state["repository_queue"]),
167 bool(state["discovery_cursor"]),
168 )
169 complete = state["discovery_complete"] and not state["repository_queue"] and not state["current_repository"]
170 return BackfillBatchResult(events=[], cursor_state=state, complete=complete)
171
172 if state["current_repository"] is None:
173 if not state["repository_queue"]:
174 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
175 state["current_repository"] = {"name": state["repository_queue"].pop(0), "cursor": None}
176
177 repository_name = state["current_repository"]["name"]
178 owner, repo_name = self._split_repository(actor, repository_name)
179 data = self.client.execute(
180 REPOSITORY_LOG_QUERY,
181 {"username": owner, "repoName": repo_name, "cursor": state["current_repository"]["cursor"]},
182 )
183 user = data.get("user") or {}
184 repository = user.get("repository") or {}
185 log_page = repository.get("log") or {}
186 commits = log_page.get("results") or []
187 next_cursor = log_page.get("cursor")
188 events = [
189 normalized
190 for commit in commits
191 if isinstance(commit, dict)
192 for normalized in [self._normalize_commit(actor=actor, repo_name=repo_name, commit=commit)]
193 if normalized is not None
194 ]
195 logger.info(
196 "git backfill actor=%s repository=%s commits=%s next_cursor=%s",
197 actor,
198 repository_name,
199 len(commits),
200 bool(next_cursor),
201 )
202 if next_cursor:
203 state["current_repository"]["cursor"] = next_cursor
204 else:
205 state["current_repository"] = None
206
207 complete = state["discovery_complete"] and not state["repository_queue"] and not state["current_repository"]
208 return BackfillBatchResult(events=events, cursor_state=None if complete else state, complete=complete)
209
120 def _discover_owned_repositories(self, actor: str) -> list[str]: 210 def _discover_owned_repositories(self, actor: str) -> list[str]:
121 owner = actor.lstrip("~") 211 owner = actor.lstrip("~")
122 repositories: list[str] = [] 212 repositories: list[str] = []
src/srht_contrib/services/todo.py +136
@@ -8,6 +8,7 @@ from typing import Any
8from srht_contrib.config import Settings 8from srht_contrib.config import Settings
9from srht_contrib.schemas import NormalizedEvent 9from srht_contrib.schemas import NormalizedEvent
10from srht_contrib.services.srht_client import SourceHutGraphQLClient 10from srht_contrib.services.srht_client import SourceHutGraphQLClient
11from srht_contrib.services.types import BackfillBatchResult
11from srht_contrib.utils.dates import ensure_utc, parse_datetime 12from srht_contrib.utils.dates import ensure_utc, parse_datetime
12 13
13 14
@@ -293,6 +294,141 @@ class TodoIngestionService:
293 tracker_events = self._fetch_from_trackers(actor=actor, since=since_dt) 294 tracker_events = self._fetch_from_trackers(actor=actor, since=since_dt)
294 return TodoPollResult(events=tracker_events, cursor=cursor_time) 295 return TodoPollResult(events=tracker_events, cursor=cursor_time)
295 296
297 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
298 state = {
299 "trackers_cursor": None,
300 "tracker_queue": [],
301 "current_tracker": None,
302 "current_ticket": None,
303 "trackers_loaded": False,
304 }
305 if cursor_state:
306 state.update(cursor_state)
307
308 if not state["trackers_loaded"]:
309 data = self.client.execute(TODO_TRACKERS_QUERY, {"cursor": state["trackers_cursor"]})
310 me = data.get("me") or {}
311 trackers_page = me.get("trackers") or {}
312 results = trackers_page.get("results") or []
313 state["trackers_cursor"] = trackers_page.get("cursor")
314 for tracker in results:
315 if not isinstance(tracker, dict):
316 continue
317 state["tracker_queue"].append(
318 {
319 "id": str(tracker.get("id")),
320 "rid": str(tracker.get("rid")),
321 "name": tracker.get("name"),
322 "tickets_cursor": None,
323 "pending_tickets": [],
324 }
325 )
326 state["trackers_loaded"] = not bool(state["trackers_cursor"])
327 logger.info(
328 "todo backfill tracker discovery actor=%s trackers=%s next_cursor=%s queue=%s",
329 actor,
330 len(results),
331 bool(state["trackers_cursor"]),
332 len(state["tracker_queue"]),
333 )
334 complete = state["trackers_loaded"] and not state["tracker_queue"]
335 return BackfillBatchResult(events=[], cursor_state=None if complete else state, complete=complete)
336
337 if state["current_tracker"] is None:
338 if not state["tracker_queue"]:
339 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
340 state["current_tracker"] = state["tracker_queue"].pop(0)
341
342 current_tracker = state["current_tracker"]
343 if state["current_ticket"] is None and not current_tracker["pending_tickets"]:
344 data = self.client.execute(
345 TODO_TRACKER_TICKETS_QUERY,
346 {"trackerRid": current_tracker["rid"], "cursor": current_tracker["tickets_cursor"]},
347 )
348 tracker = data.get("tracker") or {}
349 tickets_page = tracker.get("tickets") or {}
350 tickets = tickets_page.get("results") or []
351 current_tracker["tickets_cursor"] = tickets_page.get("cursor")
352 current_tracker["pending_tickets"].extend(
353 [{"id": int(ticket["id"]), "ref": str(ticket.get("ref") or ticket["id"])} for ticket in tickets if isinstance(ticket, dict)]
354 )
355 logger.info(
356 "todo backfill tracker=%s ticket page count=%s next_cursor=%s pending_tickets=%s",
357 current_tracker.get("name") or current_tracker.get("id"),
358 len(tickets),
359 bool(current_tracker["tickets_cursor"]),
360 len(current_tracker["pending_tickets"]),
361 )
362 if not current_tracker["pending_tickets"] and not current_tracker["tickets_cursor"]:
363 state["current_tracker"] = None
364 return BackfillBatchResult(events=[], cursor_state=state, complete=False)
365
366 if state["current_ticket"] is None:
367 if current_tracker["pending_tickets"]:
368 next_ticket = current_tracker["pending_tickets"].pop(0)
369 state["current_ticket"] = {**next_ticket, "cursor": None}
370 else:
371 state["current_tracker"] = None
372 return BackfillBatchResult(events=[], cursor_state=state, complete=False)
373
374 current_ticket = state["current_ticket"]
375 data = self.client.execute(
376 TODO_TICKET_EVENTS_QUERY,
377 {"trackerRid": current_tracker["rid"], "ticketId": current_ticket["id"], "cursor": current_ticket["cursor"]},
378 )
379 tracker = data.get("tracker") or {}
380 ticket_payload = (tracker.get("ticket") or {}) if isinstance(tracker, dict) else {}
381 event_page = ticket_payload.get("events") or {}
382 page_events = event_page.get("results") or []
383 next_cursor = event_page.get("cursor")
384 logger.info(
385 "todo backfill ticket events ref=%s tracker=%s count=%s next_cursor=%s",
386 current_ticket["ref"],
387 current_tracker.get("name") or current_tracker.get("id"),
388 len(page_events),
389 bool(next_cursor),
390 )
391
392 events: list[NormalizedEvent] = []
393 for event in page_events:
394 if not isinstance(event, dict):
395 continue
396 event["ticket"] = {
397 "id": ticket_payload.get("id", current_ticket["id"]),
398 "ref": ticket_payload.get("ref", current_ticket["ref"]),
399 "status": ticket_payload.get("status"),
400 "resolution": ticket_payload.get("resolution"),
401 "tracker": {"name": current_tracker.get("name")},
402 }
403 occurred_at = parse_datetime(event["created"])
404 for change in event.get("changes") or []:
405 if not isinstance(change, dict):
406 continue
407 normalized = _normalize_event_change(
408 settings=self.settings,
409 actor=actor,
410 event=event,
411 change=change,
412 occurred_at=occurred_at,
413 )
414 if normalized is not None:
415 events.append(normalized)
416
417 if next_cursor:
418 state["current_ticket"]["cursor"] = next_cursor
419 else:
420 state["current_ticket"] = None
421 if not current_tracker["pending_tickets"] and not current_tracker["tickets_cursor"]:
422 state["current_tracker"] = None
423
424 complete = (
425 state["trackers_loaded"]
426 and not state["tracker_queue"]
427 and state["current_tracker"] is None
428 and state["current_ticket"] is None
429 )
430 return BackfillBatchResult(events=events, cursor_state=None if complete else state, complete=complete)
431
296 def _fetch_from_activity_feed(self, actor: str, since: datetime) -> TodoPollResult: 432 def _fetch_from_activity_feed(self, actor: str, since: datetime) -> TodoPollResult:
297 events: list[NormalizedEvent] = [] 433 events: list[NormalizedEvent] = []
298 cursor: str | None = None 434 cursor: str | None = None
src/srht_contrib/services/types.py added +12
@@ -0,0 +1,12 @@
1from __future__ import annotations
2
3from dataclasses import dataclass
4
5from srht_contrib.schemas import NormalizedEvent
6
7
8@dataclass(slots=True)
9class BackfillBatchResult:
10 events: list[NormalizedEvent]
11 cursor_state: dict | None
12 complete: bool
tests/test_contributions_api.py +4
@@ -63,6 +63,9 @@ def test_public_read_registers_actor_for_lazy_indexing(client: TestClient, db_se
63 assert response.status_code == 200 63 assert response.status_code == 200
64 assert response.json()["is_indexed"] is False 64 assert response.json()["is_indexed"] is False
65 assert response.json()["indexing_state"] == "pending" 65 assert response.json()["indexing_state"] == "pending"
66 assert response.json()["is_backfilled"] is False
67 assert response.json()["backfill_state"] == "pending"
68 assert response.json()["backfill_completed_at"] is None
66 assert response.json()["last_polled_at"] is None 69 assert response.json()["last_polled_at"] is None
67 assert tracked_actor is not None 70 assert tracked_actor is not None
68 assert tracked_actor.is_active is True 71 assert tracked_actor.is_active is True
@@ -107,6 +110,7 @@ def test_contribution_stats_api(client: TestClient, db_session) -> None:
107 assert response.json()["current_streak"] == 2 110 assert response.json()["current_streak"] == 2
108 assert response.json()["is_indexed"] is True 111 assert response.json()["is_indexed"] is True
109 assert response.json()["indexing_state"] == "indexed" 112 assert response.json()["indexing_state"] == "indexed"
113 assert response.json()["is_backfilled"] is False
110 114
111 115
112def test_invalid_date_input_returns_400(client: TestClient) -> None: 116def test_invalid_date_input_returns_400(client: TestClient) -> None:
tests/test_ingestion.py +47 −1
@@ -4,10 +4,11 @@ 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, TrackedActor, TrackedRepository 7from srht_contrib.models import ServiceBackfillState, 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
11from srht_contrib.services.types import BackfillBatchResult
11 12
12 13
13class StubClient: 14class StubClient:
@@ -37,6 +38,9 @@ class RecordingTodoService:
37 events = self.events_by_call.pop(0) 38 events = self.events_by_call.pop(0)
38 return TodoPollResult(events=events, cursor="2026-03-31T00:00:00+00:00") 39 return TodoPollResult(events=events, cursor="2026-03-31T00:00:00+00:00")
39 40
41 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
42 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
43
40 44
41class EmptyGitService: 45class EmptyGitService:
42 service_name = "git" 46 service_name = "git"
@@ -54,6 +58,30 @@ class EmptyGitService:
54 def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None) -> GitPollResult: 58 def fetch_recent_events(self, actor: str, since: datetime | None = None, repositories=None) -> GitPollResult:
55 return GitPollResult(events=[], cursor="2026-03-31T00:00:00+00:00") 59 return GitPollResult(events=[], cursor="2026-03-31T00:00:00+00:00")
56 60
61 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
62 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
63
64
65class BackfillingTodoService:
66 service_name = "todo"
67
68 def fetch_recent_events(self, actor: str, since: datetime | None = None) -> TodoPollResult:
69 return TodoPollResult(events=[], cursor="2026-03-31T00:00:00+00:00")
70
71 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
72 event = NormalizedEvent(
73 service="todo",
74 event_type="ticket_created",
75 actor=actor,
76 repo_name="todo",
77 resource_id="backfill-ticket",
78 external_uid=f"todo:backfill:{actor}",
79 occurred_at=datetime(2024, 1, 1, 12, 0, tzinfo=UTC),
80 weight=1.0,
81 raw_payload_json=None,
82 )
83 return BackfillBatchResult(events=[event], cursor_state=None, complete=True)
84
57 85
58def make_settings(**overrides) -> Settings: 86def make_settings(**overrides) -> Settings:
59 values = { 87 values = {
@@ -424,3 +452,21 @@ def test_scheduled_poll_polls_known_actors_and_seeds_default_actor(db_session) -
424 assert [actor.actor for actor in tracked_actors] == ["~default", "~known"] 452 assert [actor.actor for actor in tracked_actors] == ["~default", "~known"]
425 assert all(actor.last_poll_status == "indexed" for actor in tracked_actors) 453 assert all(actor.last_poll_status == "indexed" for actor in tracked_actors)
426 assert all(actor.last_polled_at is not None for actor in tracked_actors) 454 assert all(actor.last_polled_at is not None for actor in tracked_actors)
455
456
457def test_poll_marks_backfill_complete_and_persists_service_state(db_session) -> None:
458 poller = PollerService(todo_service=BackfillingTodoService(), git_service=EmptyGitService())
459
460 inserted = poller.poll_all(db_session, "~ccleberg")
461
462 tracked_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg"))
463 service_states = db_session.scalars(
464 select(ServiceBackfillState).where(ServiceBackfillState.actor == "~ccleberg").order_by(ServiceBackfillState.service)
465 ).all()
466
467 assert inserted == 1
468 assert tracked_actor is not None
469 assert tracked_actor.backfill_status == "completed"
470 assert tracked_actor.backfill_completed_at is not None
471 assert [state.service for state in service_states] == ["git", "todo"]
472 assert all(state.status == "completed" for state in service_states)
tests/test_migrations.py +3
@@ -21,6 +21,7 @@ def test_alembic_upgrade_creates_schema(tmp_path) -> None:
21 "actor_aliases", 21 "actor_aliases",
22 "alembic_version", 22 "alembic_version",
23 "contribution_events", 23 "contribution_events",
24 "service_backfill_states",
24 "sync_states", 25 "sync_states",
25 "tracked_actors", 26 "tracked_actors",
26 "tracked_repositories", 27 "tracked_repositories",
@@ -117,6 +118,7 @@ def test_alembic_upgrade_adopts_legacy_schema(tmp_path) -> None:
117 assert "uq_tracked_repository_service_actor_name" in unique_constraints 118 assert "uq_tracked_repository_service_actor_name" in unique_constraints
118 assert actor == Settings().default_actor 119 assert actor == Settings().default_actor
119 assert "tracked_actors" in inspector.get_table_names() 120 assert "tracked_actors" in inspector.get_table_names()
121 assert "service_backfill_states" in inspector.get_table_names()
120 122
121 123
122def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) -> None: 124def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) -> None:
@@ -131,3 +133,4 @@ def test_alembic_prefers_database_url_from_environment(tmp_path, monkeypatch) ->
131 inspector = inspect(create_engine(database_url)) 133 inspector = inspect(create_engine(database_url))
132 assert "actor_aliases" in inspector.get_table_names() 134 assert "actor_aliases" in inspector.get_table_names()
133 assert "tracked_actors" in inspector.get_table_names() 135 assert "tracked_actors" in inspector.get_table_names()
136 assert "service_backfill_states" in inspector.get_table_names()