krz/hutch-stats

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

Commit 01c1b50541

01c1b50541b0dc1d42cbdaa90052f6b94ceba20c

parent: f0393a0b8d

Verified · cmc

cmc <hello@cleberg.net> · 2026-04-11 23:45 UTC

feat: prioritize recent-window backfill before full history

Layout: unified · split

API.md +16
@@ -45,6 +45,7 @@ Contribution ranges:
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- Incremental indexing and historical backfill are separate. An actor can be recently indexed without being fully backfilled yet.
48- The service prioritizes a recent visible history window first, then continues deep-history backfill afterward.
48 49
49Background polling: 50Background polling:
50 51
@@ -112,6 +113,7 @@ Behavior notes:
112- This endpoint also registers the actor for background indexing and updates the actor's `last_requested_at` timestamp. 113- This endpoint also registers the actor for background indexing and updates the actor's `last_requested_at` timestamp.
113- The response is always immediate; it does not wait for SourceHut polling to finish. 114- 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. 115- Historical backfill runs in bounded background batches and may take multiple scheduler passes to complete.
116- The recent visible window is prioritized before full-history backfill so clients can show a useful graph sooner.
115 117
116Example by year: 118Example by year:
117 119
@@ -135,6 +137,9 @@ Response `200 OK`:
135 "is_indexed": false, 137 "is_indexed": false,
136 "last_polled_at": null, 138 "last_polled_at": null,
137 "indexing_state": "pending", 139 "indexing_state": "pending",
140 "is_recent_window_backfilled": false,
141 "recent_backfill_state": "in_progress",
142 "recent_backfill_completed_at": null,
138 "is_backfilled": false, 143 "is_backfilled": false,
139 "backfill_state": "in_progress", 144 "backfill_state": "in_progress",
140 "backfill_completed_at": null, 145 "backfill_completed_at": null,
@@ -153,6 +158,9 @@ Response fields:
153- `is_indexed` boolean: whether the service has already completed at least one successful recent/incremental poll for this actor 158- `is_indexed` boolean: whether the service has already completed at least one successful recent/incremental poll for this actor
154- `last_polled_at` string or `null`: most recent successful poll time, if any 159- `last_polled_at` string or `null`: most recent successful poll time, if any
155- `indexing_state` string: one of `pending`, `indexed`, or `error` 160- `indexing_state` string: one of `pending`, `indexed`, or `error`
161- `is_recent_window_backfilled` boolean: whether the prioritized recent history window has completed backfill
162- `recent_backfill_state` string: one of `pending`, `in_progress`, `completed`, or `error`
163- `recent_backfill_completed_at` string or `null`: when recent-window backfill completed, if it has
156- `is_backfilled` boolean: whether historical backfill has completed for this actor 164- `is_backfilled` boolean: whether historical backfill has completed for this actor
157- `backfill_state` string: one of `pending`, `in_progress`, `completed`, or `error` 165- `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 166- `backfill_completed_at` string or `null`: when full historical backfill completed, if it has
@@ -169,6 +177,8 @@ Indexing state semantics:
169 177
170Backfill state semantics: 178Backfill state semantics:
171 179
180- recent-window fields:
181 - represent the prioritized visible-history window for client UX
172- `pending`: the actor has not started historical backfill yet 182- `pending`: the actor has not started historical backfill yet
173- `in_progress`: historical backfill is actively progressing in bounded background batches 183- `in_progress`: historical backfill is actively progressing in bounded background batches
174- `completed`: historical backfill has completed for all supported services 184- `completed`: historical backfill has completed for all supported services
@@ -226,6 +236,9 @@ Response `200 OK`:
226 "is_indexed": true, 236 "is_indexed": true,
227 "last_polled_at": "2026-04-11T18:05:00Z", 237 "last_polled_at": "2026-04-11T18:05:00Z",
228 "indexing_state": "indexed", 238 "indexing_state": "indexed",
239 "is_recent_window_backfilled": true,
240 "recent_backfill_state": "completed",
241 "recent_backfill_completed_at": "2026-04-11T18:02:00Z",
229 "is_backfilled": false, 242 "is_backfilled": false,
230 "backfill_state": "in_progress", 243 "backfill_state": "in_progress",
231 "backfill_completed_at": null, 244 "backfill_completed_at": null,
@@ -245,6 +258,9 @@ Response fields:
245- `is_indexed` boolean 258- `is_indexed` boolean
246- `last_polled_at` string or `null` 259- `last_polled_at` string or `null`
247- `indexing_state` string 260- `indexing_state` string
261- `is_recent_window_backfilled` boolean
262- `recent_backfill_state` string
263- `recent_backfill_completed_at` string or `null`
248- `is_backfilled` boolean 264- `is_backfilled` boolean
249- `backfill_state` string 265- `backfill_state` string
250- `backfill_completed_at` string or `null` 266- `backfill_completed_at` string or `null`
README.md +8 −1
@@ -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. It performs both recent incremental polling and bounded historical backfill. 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 recent incremental polling, prioritizes a recent visible-history window for faster UX, and then continues bounded historical backfill.
18 18
19Example use cases: 19Example use cases:
20 20
@@ -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_recent_window_backfilled": true,
212 "recent_backfill_state": "completed",
213 "recent_backfill_completed_at": "2026-04-11T18:02:00Z",
211 "is_backfilled": false, 214 "is_backfilled": false,
212 "backfill_state": "in_progress", 215 "backfill_state": "in_progress",
213 "backfill_completed_at": null, 216 "backfill_completed_at": null,
@@ -235,6 +238,9 @@ Example response:
235 "is_indexed": true, 238 "is_indexed": true,
236 "last_polled_at": "2026-04-11T18:05:00Z", 239 "last_polled_at": "2026-04-11T18:05:00Z",
237 "indexing_state": "indexed", 240 "indexing_state": "indexed",
241 "is_recent_window_backfilled": true,
242 "recent_backfill_state": "completed",
243 "recent_backfill_completed_at": "2026-04-11T18:02:00Z",
238 "is_backfilled": false, 244 "is_backfilled": false,
239 "backfill_state": "in_progress", 245 "backfill_state": "in_progress",
240 "backfill_completed_at": null, 246 "backfill_completed_at": null,
@@ -326,6 +332,7 @@ The SourceHut-specific assumptions are isolated to the service modules:
326- `git.sr.ht` polling assumes the actor's repositories are discoverable through the SourceHut GraphQL API 332- `git.sr.ht` polling assumes the actor's repositories are discoverable through the SourceHut GraphQL API
327- scheduled polling runs in-process, so it is not a distributed scheduler 333- scheduled polling runs in-process, so it is not a distributed scheduler
328- newly requested actors are indexed asynchronously, so the first public read may be empty until a scheduler or manual poll runs 334- newly requested actors are indexed asynchronously, so the first public read may be empty until a scheduler or manual poll runs
335- the recent visible-history window is prioritized first, but deep-history backfill can still take many scheduler passes for active users
329- full historical backfill can take many scheduler passes for active users because it runs in bounded batches 336- full historical backfill can take many scheduler passes for active users because it runs in bounded batches
330- alias management is config-driven; there is no alias CRUD API yet 337- alias management is config-driven; there is no alias CRUD API yet
331- current deployment model is trusted-operator V1, not a public multi-tenant service 338- current deployment model is trusted-operator V1, not a public multi-tenant service
alembic/versions/20260411_0004_recent_backfill_scope.py added +89
@@ -0,0 +1,89 @@
1"""recent backfill scope and actor fields"""
2
3from __future__ import annotations
4
5from alembic import op
6import sqlalchemy as sa
7from sqlalchemy import inspect
8
9
10revision = "20260411_0004"
11down_revision = "20260411_0003"
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 _service_backfill_needs_upgrade() -> bool:
25 if "service_backfill_states" not in _table_names():
26 return False
27 columns = _column_names("service_backfill_states")
28 if "scope" not in columns:
29 return True
30 unique_constraints = {
31 constraint["name"]
32 for constraint in inspect(op.get_bind()).get_unique_constraints("service_backfill_states")
33 }
34 return "uq_service_backfill_state_actor_service_scope" not in unique_constraints
35
36
37def _upgrade_service_backfill_states() -> None:
38 op.execute(
39 sa.text(
40 """
41 CREATE TABLE service_backfill_states__alembic_new (
42 id INTEGER NOT NULL PRIMARY KEY,
43 actor VARCHAR(255) NOT NULL,
44 service VARCHAR(32) NOT NULL,
45 scope VARCHAR(16) NOT NULL,
46 cursor_json JSON,
47 status VARCHAR(32) NOT NULL,
48 started_at DATETIME,
49 completed_at DATETIME,
50 last_error TEXT,
51 updated_at DATETIME NOT NULL,
52 CONSTRAINT uq_service_backfill_state_actor_service_scope UNIQUE (actor, service, scope)
53 )
54 """
55 )
56 )
57 op.execute(
58 sa.text(
59 """
60 INSERT INTO service_backfill_states__alembic_new
61 (id, actor, service, scope, cursor_json, status, started_at, completed_at, last_error, updated_at)
62 SELECT id, actor, service, 'full', cursor_json, status, started_at, completed_at, last_error, updated_at
63 FROM service_backfill_states
64 """
65 )
66 )
67 op.execute(sa.text("DROP TABLE service_backfill_states"))
68 op.execute(sa.text("ALTER TABLE service_backfill_states__alembic_new RENAME TO service_backfill_states"))
69
70
71def upgrade() -> None:
72 if "tracked_actors" in _table_names():
73 columns = _column_names("tracked_actors")
74 if "recent_backfill_status" not in columns:
75 op.add_column("tracked_actors", sa.Column("recent_backfill_status", sa.String(length=32), nullable=True))
76 op.execute(sa.text("UPDATE tracked_actors SET recent_backfill_status = 'pending' WHERE recent_backfill_status IS NULL"))
77 if "recent_backfill_started_at" not in columns:
78 op.add_column("tracked_actors", sa.Column("recent_backfill_started_at", sa.DateTime(timezone=True), nullable=True))
79 if "recent_backfill_completed_at" not in columns:
80 op.add_column("tracked_actors", sa.Column("recent_backfill_completed_at", sa.DateTime(timezone=True), nullable=True))
81 if "last_recent_backfill_error" not in columns:
82 op.add_column("tracked_actors", sa.Column("last_recent_backfill_error", sa.Text(), nullable=True))
83
84 if _service_backfill_needs_upgrade():
85 _upgrade_service_backfill_states()
86
87
88def downgrade() -> None:
89 pass
src/srht_contrib/jobs/poller.py +128 −45
@@ -18,6 +18,9 @@ from srht_contrib.utils.repositories import canonicalize_repository_name
18 18
19logger = logging.getLogger(__name__) 19logger = logging.getLogger(__name__)
20SYNC_OVERLAP = timedelta(hours=24) 20SYNC_OVERLAP = timedelta(hours=24)
21RECENT_BACKFILL_DAYS = 365
22RECENT_BACKFILL_BATCHES_PER_SERVICE = 5
23FULL_BACKFILL_BATCHES_PER_SERVICE = 1
21 24
22 25
23class PollerService: 26class PollerService:
@@ -63,7 +66,12 @@ class PollerService:
63 def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor: 66 def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor:
64 tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) 67 tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor))
65 if tracked_actor is None: 68 if tracked_actor is None:
66 tracked_actor = TrackedActor(actor=actor, is_active=True, backfill_status="pending") 69 tracked_actor = TrackedActor(
70 actor=actor,
71 is_active=True,
72 recent_backfill_status="pending",
73 backfill_status="pending",
74 )
67 db.add(tracked_actor) 75 db.add(tracked_actor)
68 76
69 tracked_actor.is_active = True 77 tracked_actor.is_active = True
@@ -177,32 +185,98 @@ class PollerService:
177 185
178 def _run_backfill_batches(self, db: Session, actor: str) -> int: 186 def _run_backfill_batches(self, db: Session, actor: str) -> int:
179 tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) 187 tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
180 if tracked_actor.backfill_status == "completed": 188 if tracked_actor.recent_backfill_status == "completed" and tracked_actor.backfill_status == "completed":
181 return 0 189 return 0
182 190
183 if tracked_actor.backfill_started_at is None: 191 now = datetime.now(tz=UTC)
184 tracked_actor.backfill_started_at = datetime.now(tz=UTC) 192 recent_since = now - timedelta(days=RECENT_BACKFILL_DAYS)
185 tracked_actor.backfill_status = "in_progress"
186 tracked_actor.last_backfill_error = None
187 db.add(tracked_actor) 193 db.add(tracked_actor)
188 db.flush() 194 db.flush()
189 total_inserted = 0 195 total_inserted = 0
190 196
191 services = [ 197 recent_services = [
192 (self.todo_service.service_name, self.todo_service.fetch_backfill_batch), 198 (self.todo_service.service_name, self.todo_service.fetch_recent_backfill_batch),
193 (self.git_service.service_name, self.git_service.fetch_backfill_batch), 199 (self.git_service.service_name, self.git_service.fetch_recent_backfill_batch),
194 ] 200 ]
195 all_complete = True 201 if tracked_actor.recent_backfill_status != "completed":
202 if tracked_actor.recent_backfill_started_at is None:
203 tracked_actor.recent_backfill_started_at = now
204 tracked_actor.recent_backfill_status = "in_progress"
205 tracked_actor.last_recent_backfill_error = None
206 total_inserted += self._run_backfill_scope(
207 db,
208 actor=actor,
209 scope="recent",
210 services=recent_services,
211 batches_per_service=RECENT_BACKFILL_BATCHES_PER_SERVICE,
212 since=recent_since,
213 )
214 recent_statuses = db.scalars(
215 select(ServiceBackfillState.status)
216 .where(ServiceBackfillState.actor == actor)
217 .where(ServiceBackfillState.scope == "recent")
218 ).all()
219 if recent_statuses and all(status == "completed" for status in recent_statuses):
220 tracked_actor.recent_backfill_status = "completed"
221 tracked_actor.recent_backfill_completed_at = datetime.now(tz=UTC)
222 tracked_actor.last_recent_backfill_error = None
223 db.add(tracked_actor)
224 db.flush()
225
226 if tracked_actor.recent_backfill_status == "completed" and tracked_actor.backfill_status != "completed":
227 if tracked_actor.backfill_started_at is None:
228 tracked_actor.backfill_started_at = now
229 tracked_actor.backfill_status = "in_progress"
230 tracked_actor.last_backfill_error = None
231 total_inserted += self._run_backfill_scope(
232 db,
233 actor=actor,
234 scope="full",
235 services=[
236 (self.todo_service.service_name, self.todo_service.fetch_backfill_batch),
237 (self.git_service.service_name, self.git_service.fetch_backfill_batch),
238 ],
239 batches_per_service=FULL_BACKFILL_BATCHES_PER_SERVICE,
240 since=None,
241 )
242 full_statuses = db.scalars(
243 select(ServiceBackfillState.status)
244 .where(ServiceBackfillState.actor == actor)
245 .where(ServiceBackfillState.scope == "full")
246 ).all()
247 if full_statuses and all(status == "completed" for status in full_statuses):
248 tracked_actor.backfill_status = "completed"
249 tracked_actor.backfill_completed_at = datetime.now(tz=UTC)
250 tracked_actor.last_backfill_error = None
251 db.add(tracked_actor)
252 db.flush()
253
254 return total_inserted
255
256 def _run_backfill_scope(
257 self,
258 db: Session,
259 *,
260 actor: str,
261 scope: str,
262 services,
263 batches_per_service: int,
264 since: datetime | None,
265 ) -> int:
266 tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
267 total_inserted = 0
196 for service_name, fetcher in services: 268 for service_name, fetcher in services:
197 state = db.scalar( 269 state = db.scalar(
198 select(ServiceBackfillState) 270 select(ServiceBackfillState)
199 .where(ServiceBackfillState.actor == actor) 271 .where(ServiceBackfillState.actor == actor)
200 .where(ServiceBackfillState.service == service_name) 272 .where(ServiceBackfillState.service == service_name)
273 .where(ServiceBackfillState.scope == scope)
201 ) 274 )
202 if state is None: 275 if state is None:
203 state = ServiceBackfillState( 276 state = ServiceBackfillState(
204 actor=actor, 277 actor=actor,
205 service=service_name, 278 service=service_name,
279 scope=scope,
206 cursor_json=None, 280 cursor_json=None,
207 status="pending", 281 status="pending",
208 started_at=None, 282 started_at=None,
@@ -216,47 +290,56 @@ class PollerService:
216 if state.status == "completed": 290 if state.status == "completed":
217 continue 291 continue
218 292
219 all_complete = False
220 if state.started_at is None: 293 if state.started_at is None:
221 state.started_at = datetime.now(tz=UTC) 294 state.started_at = datetime.now(tz=UTC)
222 state.status = "in_progress" 295 state.status = "in_progress"
223 state.updated_at = datetime.now(tz=UTC) 296 state.updated_at = datetime.now(tz=UTC)
224 try: 297
225 result = fetcher(actor=actor, cursor_state=state.cursor_json) 298 for _ in range(batches_per_service):
226 inserted = self._insert_events(db, result.events) 299 try:
227 total_inserted += inserted 300 if since is None:
228 state.cursor_json = copy.deepcopy(result.cursor_state) 301 result = fetcher(actor=actor, cursor_state=state.cursor_json)
229 state.last_error = None 302 else:
230 state.updated_at = datetime.now(tz=UTC) 303 result = fetcher(actor=actor, cursor_state=state.cursor_json, since=since)
231 if result.complete: 304 inserted = self._insert_events(db, result.events)
232 state.status = "completed" 305 total_inserted += inserted
233 state.completed_at = datetime.now(tz=UTC) 306 state.cursor_json = copy.deepcopy(result.cursor_state)
234 logger.info("Backfill complete for service=%s actor=%s inserted=%s", service_name, actor, inserted) 307 state.last_error = None
235 else: 308 state.updated_at = datetime.now(tz=UTC)
236 logger.info("Backfill batch complete for service=%s actor=%s inserted=%s", service_name, actor, inserted) 309 if result.complete:
237 except Exception as exc: 310 state.status = "completed"
238 state.status = "error" 311 state.completed_at = datetime.now(tz=UTC)
239 state.last_error = str(exc) 312 logger.info(
240 state.updated_at = datetime.now(tz=UTC) 313 "Backfill complete for scope=%s service=%s actor=%s inserted=%s",
241 tracked_actor.backfill_status = "error" 314 scope,
242 tracked_actor.last_backfill_error = str(exc) 315 service_name,
243 db.add(state) 316 actor,
244 db.add(tracked_actor) 317 inserted,
245 db.flush() 318 )
246 raise 319 break
320 logger.info(
321 "Backfill batch complete for scope=%s service=%s actor=%s inserted=%s",
322 scope,
323 service_name,
324 actor,
325 inserted,
326 )
327 except Exception as exc:
328 state.status = "error"
329 state.last_error = str(exc)
330 state.updated_at = datetime.now(tz=UTC)
331 if scope == "recent":
332 tracked_actor.recent_backfill_status = "error"
333 tracked_actor.last_recent_backfill_error = str(exc)
334 else:
335 tracked_actor.backfill_status = "error"
336 tracked_actor.last_backfill_error = str(exc)
337 db.add(state)
338 db.add(tracked_actor)
339 db.flush()
340 raise
247 341
248 db.add(state) 342 db.add(state)
249 db.flush() 343 db.flush()
250 344
251 completed = db.scalars(
252 select(ServiceBackfillState.status).where(ServiceBackfillState.actor == actor)
253 ).all()
254 if completed and all(status == "completed" for status in completed):
255 tracked_actor.backfill_status = "completed"
256 tracked_actor.backfill_completed_at = datetime.now(tz=UTC)
257 tracked_actor.last_backfill_error = None
258 elif tracked_actor.backfill_status != "error":
259 tracked_actor.backfill_status = "in_progress"
260 db.add(tracked_actor)
261 db.flush()
262 return total_inserted 345 return total_inserted
src/srht_contrib/models.py +6 −1
@@ -71,6 +71,10 @@ 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 recent_backfill_status: Mapped[str] = mapped_column(String(32), nullable=False, default="pending")
75 recent_backfill_started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
76 recent_backfill_completed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
77 last_recent_backfill_error: Mapped[str | None] = mapped_column(Text, nullable=True)
74 backfill_status: Mapped[str] = mapped_column(String(32), nullable=False, default="pending") 78 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) 79 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) 80 backfill_completed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
@@ -79,11 +83,12 @@ class TrackedActor(Base):
79 83
80class ServiceBackfillState(Base): 84class ServiceBackfillState(Base):
81 __tablename__ = "service_backfill_states" 85 __tablename__ = "service_backfill_states"
82 __table_args__ = (UniqueConstraint("actor", "service", name="uq_service_backfill_state_actor_service"),) 86 __table_args__ = (UniqueConstraint("actor", "service", "scope", name="uq_service_backfill_state_actor_service_scope"),)
83 87
84 id: Mapped[int] = mapped_column(Integer, primary_key=True) 88 id: Mapped[int] = mapped_column(Integer, primary_key=True)
85 actor: Mapped[str] = mapped_column(String(255), nullable=False) 89 actor: Mapped[str] = mapped_column(String(255), nullable=False)
86 service: Mapped[str] = mapped_column(String(32), nullable=False) 90 service: Mapped[str] = mapped_column(String(32), nullable=False)
91 scope: Mapped[str] = mapped_column(String(16), nullable=False, default="full")
87 cursor_json: Mapped[dict | None] = mapped_column(JSON, nullable=True) 92 cursor_json: Mapped[dict | None] = mapped_column(JSON, nullable=True)
88 status: Mapped[str] = mapped_column(String(32), nullable=False, default="pending") 93 status: Mapped[str] = mapped_column(String(32), nullable=False, default="pending")
89 started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) 94 started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
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_recent_window_backfilled: bool = False
32 recent_backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending"
33 recent_backfill_completed_at: datetime | None = None
31 is_backfilled: bool = False 34 is_backfilled: bool = False
32 backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending" 35 backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending"
33 backfill_completed_at: datetime | None = None 36 backfill_completed_at: datetime | None = None
src/srht_contrib/services/aggregator.py +9
@@ -62,6 +62,12 @@ class ContributionAggregator:
62 is_indexed=calendar.is_indexed, 62 is_indexed=calendar.is_indexed,
63 last_polled_at=calendar.last_polled_at, 63 last_polled_at=calendar.last_polled_at,
64 indexing_state=calendar.indexing_state, 64 indexing_state=calendar.indexing_state,
65 is_recent_window_backfilled=calendar.is_recent_window_backfilled,
66 recent_backfill_state=calendar.recent_backfill_state,
67 recent_backfill_completed_at=calendar.recent_backfill_completed_at,
68 is_backfilled=calendar.is_backfilled,
69 backfill_state=calendar.backfill_state,
70 backfill_completed_at=calendar.backfill_completed_at,
65 ) 71 )
66 72
67 def _index_metadata(self, db: Session, actor: str) -> ContributionIndexMetadata: 73 def _index_metadata(self, db: Session, actor: str) -> ContributionIndexMetadata:
@@ -81,6 +87,9 @@ class ContributionAggregator:
81 is_indexed=is_indexed, 87 is_indexed=is_indexed,
82 last_polled_at=tracked_actor.last_polled_at if tracked_actor is not None else None, 88 last_polled_at=tracked_actor.last_polled_at if tracked_actor is not None else None,
83 indexing_state=indexing_state, 89 indexing_state=indexing_state,
90 is_recent_window_backfilled=(tracked_actor.recent_backfill_status == "completed") if tracked_actor is not None else False,
91 recent_backfill_state=(tracked_actor.recent_backfill_status if tracked_actor is not None else "pending"),
92 recent_backfill_completed_at=tracked_actor.recent_backfill_completed_at if tracked_actor is not None else None,
84 is_backfilled=(tracked_actor.backfill_status == "completed") if tracked_actor is not None else False, 93 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"), 94 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, 95 backfill_completed_at=tracked_actor.backfill_completed_at if tracked_actor is not None else None,
src/srht_contrib/services/git.py +31 −8
@@ -120,6 +120,24 @@ class GitIngestionService:
120 return repositories 120 return repositories
121 121
122 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: 122 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
123 return self._fetch_backfill_batch(actor, cursor_state, since=None)
124
125 def fetch_recent_backfill_batch(
126 self,
127 actor: str,
128 cursor_state: dict | None = None,
129 *,
130 since: datetime,
131 ) -> BackfillBatchResult:
132 return self._fetch_backfill_batch(actor, cursor_state, since=since)
133
134 def _fetch_backfill_batch(
135 self,
136 actor: str,
137 cursor_state: dict | None,
138 *,
139 since: datetime | None,
140 ) -> BackfillBatchResult:
123 state = { 141 state = {
124 "discovery_cursor": None, 142 "discovery_cursor": None,
125 "discovery_complete": False, 143 "discovery_complete": False,
@@ -186,13 +204,18 @@ class GitIngestionService:
186 log_page = repository.get("log") or {} 204 log_page = repository.get("log") or {}
187 commits = log_page.get("results") or [] 205 commits = log_page.get("results") or []
188 next_cursor = log_page.get("cursor") 206 next_cursor = log_page.get("cursor")
189 events = [ 207 events: list[NormalizedEvent] = []
190 normalized 208 stop_repository = False
191 for commit in commits 209 for commit in commits:
192 if isinstance(commit, dict) 210 if not isinstance(commit, dict):
193 for normalized in [self._normalize_commit(actor=actor, repo_name=repo_name, commit=commit)] 211 continue
194 if normalized is not None 212 commit_time = parse_datetime((commit.get("author") or {}).get("time"))
195 ] 213 if since is not None and commit_time < since:
214 stop_repository = True
215 break
216 normalized = self._normalize_commit(actor=actor, repo_name=repo_name, commit=commit)
217 if normalized is not None:
218 events.append(normalized)
196 logger.info( 219 logger.info(
197 "git backfill actor=%s repository=%s commits=%s next_cursor=%s", 220 "git backfill actor=%s repository=%s commits=%s next_cursor=%s",
198 actor, 221 actor,
@@ -200,7 +223,7 @@ class GitIngestionService:
200 len(commits), 223 len(commits),
201 bool(next_cursor), 224 bool(next_cursor),
202 ) 225 )
203 if next_cursor: 226 if next_cursor and not stop_repository:
204 state["current_repository"]["cursor"] = next_cursor 227 state["current_repository"]["cursor"] = next_cursor
205 else: 228 else:
206 state["current_repository"] = None 229 state["current_repository"] = None
src/srht_contrib/services/todo.py +37 −5
@@ -296,6 +296,24 @@ class TodoIngestionService:
296 return TodoPollResult(events=tracker_events, cursor=cursor_time) 296 return TodoPollResult(events=tracker_events, cursor=cursor_time)
297 297
298 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: 298 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
299 return self._fetch_backfill_batch(actor, cursor_state, since=None)
300
301 def fetch_recent_backfill_batch(
302 self,
303 actor: str,
304 cursor_state: dict | None = None,
305 *,
306 since: datetime,
307 ) -> BackfillBatchResult:
308 return self._fetch_backfill_batch(actor, cursor_state, since=since)
309
310 def _fetch_backfill_batch(
311 self,
312 actor: str,
313 cursor_state: dict | None,
314 *,
315 since: datetime | None,
316 ) -> BackfillBatchResult:
299 state = { 317 state = {
300 "trackers_cursor": None, 318 "trackers_cursor": None,
301 "tracker_queue": [], 319 "tracker_queue": [],
@@ -349,10 +367,20 @@ class TodoIngestionService:
349 tracker = data.get("tracker") or {} 367 tracker = data.get("tracker") or {}
350 tickets_page = tracker.get("tickets") or {} 368 tickets_page = tracker.get("tickets") or {}
351 tickets = tickets_page.get("results") or [] 369 tickets = tickets_page.get("results") or []
352 current_tracker["tickets_cursor"] = tickets_page.get("cursor") 370 next_tickets_cursor = tickets_page.get("cursor")
353 current_tracker["pending_tickets"].extend( 371 stop_tracker_paging = False
354 [{"id": int(ticket["id"]), "ref": str(ticket.get("ref") or ticket["id"])} for ticket in tickets if isinstance(ticket, dict)] 372 for ticket in tickets:
355 ) 373 if not isinstance(ticket, dict):
374 continue
375 if since is not None:
376 updated = parse_datetime(ticket["updated"])
377 if updated < since:
378 stop_tracker_paging = True
379 continue
380 current_tracker["pending_tickets"].append(
381 {"id": int(ticket["id"]), "ref": str(ticket.get("ref") or ticket["id"])}
382 )
383 current_tracker["tickets_cursor"] = None if stop_tracker_paging else next_tickets_cursor
356 logger.info( 384 logger.info(
357 "todo backfill tracker=%s ticket page count=%s next_cursor=%s pending_tickets=%s", 385 "todo backfill tracker=%s ticket page count=%s next_cursor=%s pending_tickets=%s",
358 current_tracker.get("name") or current_tracker.get("id"), 386 current_tracker.get("name") or current_tracker.get("id"),
@@ -391,6 +419,7 @@ class TodoIngestionService:
391 ) 419 )
392 420
393 events: list[NormalizedEvent] = [] 421 events: list[NormalizedEvent] = []
422 stop_ticket_paging = False
394 for event in page_events: 423 for event in page_events:
395 if not isinstance(event, dict): 424 if not isinstance(event, dict):
396 continue 425 continue
@@ -402,6 +431,9 @@ class TodoIngestionService:
402 "tracker": {"name": current_tracker.get("name")}, 431 "tracker": {"name": current_tracker.get("name")},
403 } 432 }
404 occurred_at = parse_datetime(event["created"]) 433 occurred_at = parse_datetime(event["created"])
434 if since is not None and occurred_at < since:
435 stop_ticket_paging = True
436 continue
405 for change in event.get("changes") or []: 437 for change in event.get("changes") or []:
406 if not isinstance(change, dict): 438 if not isinstance(change, dict):
407 continue 439 continue
@@ -415,7 +447,7 @@ class TodoIngestionService:
415 if normalized is not None: 447 if normalized is not None:
416 events.append(normalized) 448 events.append(normalized)
417 449
418 if next_cursor: 450 if next_cursor and not stop_ticket_paging:
419 state["current_ticket"]["cursor"] = next_cursor 451 state["current_ticket"]["cursor"] = next_cursor
420 else: 452 else:
421 state["current_ticket"] = None 453 state["current_ticket"] = None
tests/test_contributions_api.py +7
@@ -48,6 +48,8 @@ def test_contributions_api_returns_zero_filled_range(client: TestClient, db_sess
48 assert response.status_code == 200 48 assert response.status_code == 200
49 assert response.json()["is_indexed"] is True 49 assert response.json()["is_indexed"] is True
50 assert response.json()["indexing_state"] == "indexed" 50 assert response.json()["indexing_state"] == "indexed"
51 assert response.json()["is_recent_window_backfilled"] is False
52 assert response.json()["recent_backfill_state"] == "pending"
51 assert response.json()["days"] == [ 53 assert response.json()["days"] == [
52 {"date": "2026-03-28", "count": 0, "score": 0.0}, 54 {"date": "2026-03-28", "count": 0, "score": 0.0},
53 {"date": "2026-03-29", "count": 0, "score": 0.0}, 55 {"date": "2026-03-29", "count": 0, "score": 0.0},
@@ -63,6 +65,9 @@ def test_public_read_registers_actor_for_lazy_indexing(client: TestClient, db_se
63 assert response.status_code == 200 65 assert response.status_code == 200
64 assert response.json()["is_indexed"] is False 66 assert response.json()["is_indexed"] is False
65 assert response.json()["indexing_state"] == "pending" 67 assert response.json()["indexing_state"] == "pending"
68 assert response.json()["is_recent_window_backfilled"] is False
69 assert response.json()["recent_backfill_state"] == "pending"
70 assert response.json()["recent_backfill_completed_at"] is None
66 assert response.json()["is_backfilled"] is False 71 assert response.json()["is_backfilled"] is False
67 assert response.json()["backfill_state"] == "pending" 72 assert response.json()["backfill_state"] == "pending"
68 assert response.json()["backfill_completed_at"] is None 73 assert response.json()["backfill_completed_at"] is None
@@ -110,6 +115,8 @@ def test_contribution_stats_api(client: TestClient, db_session) -> None:
110 assert response.json()["current_streak"] == 2 115 assert response.json()["current_streak"] == 2
111 assert response.json()["is_indexed"] is True 116 assert response.json()["is_indexed"] is True
112 assert response.json()["indexing_state"] == "indexed" 117 assert response.json()["indexing_state"] == "indexed"
118 assert response.json()["is_recent_window_backfilled"] is False
119 assert response.json()["recent_backfill_state"] == "pending"
113 assert response.json()["is_backfilled"] is False 120 assert response.json()["is_backfilled"] is False
114 121
115 122
tests/test_ingestion.py +57 −4
@@ -41,6 +41,15 @@ class RecordingTodoService:
41 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: 41 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
42 return BackfillBatchResult(events=[], cursor_state=None, complete=True) 42 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
43 43
44 def fetch_recent_backfill_batch(
45 self,
46 actor: str,
47 cursor_state: dict | None = None,
48 *,
49 since: datetime,
50 ) -> BackfillBatchResult:
51 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
52
44 53
45class EmptyGitService: 54class EmptyGitService:
46 service_name = "git" 55 service_name = "git"
@@ -61,6 +70,15 @@ class EmptyGitService:
61 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: 70 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
62 return BackfillBatchResult(events=[], cursor_state=None, complete=True) 71 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
63 72
73 def fetch_recent_backfill_batch(
74 self,
75 actor: str,
76 cursor_state: dict | None = None,
77 *,
78 since: datetime,
79 ) -> BackfillBatchResult:
80 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
81
64 82
65class BackfillingTodoService: 83class BackfillingTodoService:
66 service_name = "todo" 84 service_name = "todo"
@@ -82,6 +100,15 @@ class BackfillingTodoService:
82 ) 100 )
83 return BackfillBatchResult(events=[event], cursor_state=None, complete=True) 101 return BackfillBatchResult(events=[event], cursor_state=None, complete=True)
84 102
103 def fetch_recent_backfill_batch(
104 self,
105 actor: str,
106 cursor_state: dict | None = None,
107 *,
108 since: datetime,
109 ) -> BackfillBatchResult:
110 return self.fetch_backfill_batch(actor, cursor_state)
111
85 112
86class QueueShrinkingTodoService: 113class QueueShrinkingTodoService:
87 service_name = "todo" 114 service_name = "todo"
@@ -100,6 +127,15 @@ class QueueShrinkingTodoService:
100 state["tracker_queue"].pop(0) 127 state["tracker_queue"].pop(0)
101 return BackfillBatchResult(events=[], cursor_state=state, complete=False) 128 return BackfillBatchResult(events=[], cursor_state=state, complete=False)
102 129
130 def fetch_recent_backfill_batch(
131 self,
132 actor: str,
133 cursor_state: dict | None = None,
134 *,
135 since: datetime,
136 ) -> BackfillBatchResult:
137 return self.fetch_backfill_batch(actor, cursor_state)
138
103 139
104class QueueShrinkingGitService: 140class QueueShrinkingGitService:
105 service_name = "git" 141 service_name = "git"
@@ -128,6 +164,15 @@ class QueueShrinkingGitService:
128 state["repository_queue"].pop(0) 164 state["repository_queue"].pop(0)
129 return BackfillBatchResult(events=[], cursor_state=state, complete=False) 165 return BackfillBatchResult(events=[], cursor_state=state, complete=False)
130 166
167 def fetch_recent_backfill_batch(
168 self,
169 actor: str,
170 cursor_state: dict | None = None,
171 *,
172 since: datetime,
173 ) -> BackfillBatchResult:
174 return self.fetch_backfill_batch(actor, cursor_state)
175
131 176
132def make_settings(**overrides) -> Settings: 177def make_settings(**overrides) -> Settings:
133 values = { 178 values = {
@@ -507,14 +552,18 @@ def test_poll_marks_backfill_complete_and_persists_service_state(db_session) ->
507 552
508 tracked_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg")) 553 tracked_actor = db_session.scalar(select(TrackedActor).where(TrackedActor.actor == "~ccleberg"))
509 service_states = db_session.scalars( 554 service_states = db_session.scalars(
510 select(ServiceBackfillState).where(ServiceBackfillState.actor == "~ccleberg").order_by(ServiceBackfillState.service) 555 select(ServiceBackfillState)
556 .where(ServiceBackfillState.actor == "~ccleberg")
557 .order_by(ServiceBackfillState.scope, ServiceBackfillState.service)
511 ).all() 558 ).all()
512 559
513 assert inserted == 1 560 assert inserted == 1
514 assert tracked_actor is not None 561 assert tracked_actor is not None
562 assert tracked_actor.recent_backfill_status == "completed"
563 assert tracked_actor.recent_backfill_completed_at is not None
515 assert tracked_actor.backfill_status == "completed" 564 assert tracked_actor.backfill_status == "completed"
516 assert tracked_actor.backfill_completed_at is not None 565 assert tracked_actor.backfill_completed_at is not None
517 assert [state.service for state in service_states] == ["git", "todo"] 566 assert [f"{state.scope}:{state.service}" for state in service_states] == ["full:git", "full:todo", "recent:git", "recent:todo"]
518 assert all(state.status == "completed" for state in service_states) 567 assert all(state.status == "completed" for state in service_states)
519 568
520 569
@@ -525,7 +574,9 @@ def test_backfill_cursor_state_shrinks_across_repeated_polls(db_session) -> None
525 first_states = { 574 first_states = {
526 state.service: state.cursor_json 575 state.service: state.cursor_json
527 for state in db_session.scalars( 576 for state in db_session.scalars(
528 select(ServiceBackfillState).where(ServiceBackfillState.actor == "~ccleberg") 577 select(ServiceBackfillState)
578 .where(ServiceBackfillState.actor == "~ccleberg")
579 .where(ServiceBackfillState.scope == "full")
529 ).all() 580 ).all()
530 } 581 }
531 582
@@ -533,7 +584,9 @@ def test_backfill_cursor_state_shrinks_across_repeated_polls(db_session) -> None
533 second_states = { 584 second_states = {
534 state.service: state.cursor_json 585 state.service: state.cursor_json
535 for state in db_session.scalars( 586 for state in db_session.scalars(
536 select(ServiceBackfillState).where(ServiceBackfillState.actor == "~ccleberg") 587 select(ServiceBackfillState)
588 .where(ServiceBackfillState.actor == "~ccleberg")
589 .where(ServiceBackfillState.scope == "full")
537 ).all() 590 ).all()
538 } 591 }
539 592