krz/hutch-stats
Server-side utility for calculating contributions for sourcehut users.
clone: git clone https://gitbay.org/krz/hutch-stats.git
main: src/srht_contrib/jobs/poller.py · raw
1from __future__ import annotations
2
3import copy
4import logging
5from datetime import UTC, datetime, timedelta
6
7from sqlalchemy import or_, select
8from sqlalchemy.exc import IntegrityError
9from sqlalchemy.orm import Session
10
11from srht_contrib.config import Settings
12from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncState, TrackedActor, TrackedRepository
13from srht_contrib.schemas import NormalizedEvent
14from srht_contrib.services.git import GitIngestionService
15from srht_contrib.services.srht_client import SourceHutClientError
16from srht_contrib.services.todo import TodoIngestionService
17from srht_contrib.utils.retention import RETENTION_DAYS, prune_contribution_events
18from srht_contrib.utils.repositories import canonicalize_repository_name
19
20
21logger = logging.getLogger(__name__)
22RECENT_BACKFILL_BATCHES_PER_SERVICE = 5
23
24
25class PollerService:
26 def __init__(
27 self,
28 todo_service: TodoIngestionService,
29 git_service: GitIngestionService,
30 settings: Settings,
31 ) -> None:
32 self.todo_service = todo_service
33 self.git_service = git_service
34 self.settings = settings
35 self._sync_overlap = timedelta(hours=settings.sync_overlap_hours)
36
37 def poll_all(self, db: Session, actor: str, *, run_backfill: bool = True) -> int:
38 self.track_actor_request(db, actor, update_last_requested=False)
39 try:
40 inserted = self._poll_actor(db, actor)
41 except Exception as exc:
42 db.rollback()
43 self._update_tracked_actor_poll_state(db, actor, status="error", error=str(exc))
44 db.commit()
45 raise
46
47 self._update_tracked_actor_poll_state(db, actor, status="indexed", error=None)
48 if run_backfill:
49 inserted += self._run_backfill_batches(db, actor)
50 db.commit()
51 return inserted
52
53 def poll_tracked_actors(self, db: Session, default_actor: str | None = None) -> dict[str, int]:
54 if default_actor:
55 self.track_actor_request(db, default_actor, update_last_requested=False)
56 db.commit()
57
58 results: dict[str, int] = {}
59 actors = db.scalars(
60 select(TrackedActor)
61 .where(TrackedActor.is_active.is_(True))
62 .where(
63 or_(
64 TrackedActor.next_poll_after.is_(None),
65 TrackedActor.next_poll_after <= datetime.now(tz=UTC),
66 )
67 )
68 .order_by(
69 TrackedActor.priority_boosted_at.is_not(None).desc(),
70 TrackedActor.priority_boosted_at.desc(),
71 TrackedActor.next_poll_after.is_(None).desc(),
72 TrackedActor.next_poll_after,
73 TrackedActor.queued_for_discovery_at,
74 TrackedActor.actor,
75 )
76 .limit(self.settings.discovery_batch_size)
77 ).all()
78 for tracked_actor in actors:
79 actor = tracked_actor.actor
80 run_backfill_only = (
81 tracked_actor.last_polled_at is not None
82 and tracked_actor.recent_backfill_status != "completed"
83 )
84 claimed_actor = self.track_actor_request(db, actor, update_last_requested=False)
85 claimed_actor.discovery_state = "in_progress"
86 claimed_actor.last_claimed_at = datetime.now(tz=UTC)
87 claimed_actor.poll_attempts += 1
88 db.add(claimed_actor)
89 db.commit()
90 try:
91 if run_backfill_only:
92 results[actor] = self.poll_recent_backfill(db, actor)
93 else:
94 results[actor] = self.poll_all(db, actor, run_backfill=False)
95 self._schedule_pending_backfill_or_repoll(db, actor)
96 db.commit()
97 except SourceHutClientError:
98 logger.exception("Scheduled poll failed for actor=%s", actor)
99 except Exception:
100 logger.exception("Unexpected scheduled poll failure for actor=%s", actor)
101 deleted = self.prune_old_events(db)
102 if deleted:
103 logger.info("Pruned %s contribution events older than %s days", deleted, RETENTION_DAYS)
104 return results
105
106 def track_actor_request(
107 self,
108 db: Session,
109 actor: str,
110 *,
111 update_last_requested: bool = True,
112 prioritize: bool = False,
113 ) -> TrackedActor:
114 tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor))
115 now = datetime.now(tz=UTC)
116 if tracked_actor is None:
117 tracked_actor = TrackedActor(
118 actor=actor,
119 is_active=True,
120 discovery_state="queued",
121 queued_for_discovery_at=now,
122 next_poll_after=now,
123 recent_backfill_status="pending",
124 )
125 db.add(tracked_actor)
126
127 tracked_actor.is_active = True
128 if tracked_actor.queued_for_discovery_at is None:
129 tracked_actor.queued_for_discovery_at = now
130 if tracked_actor.next_poll_after is None:
131 tracked_actor.next_poll_after = now
132 if update_last_requested:
133 tracked_actor.last_requested_at = now
134 if prioritize:
135 tracked_actor.priority_boosted_at = now
136 tracked_actor.next_poll_after = now
137 db.flush()
138 return tracked_actor
139
140 def _poll_actor(self, db: Session, actor: str) -> int:
141 inserted = 0
142 inserted += self._poll_service(db, actor, self.todo_service.service_name, self.todo_service.fetch_recent_events)
143 self._sync_tracked_repositories(db, actor)
144 inserted += self._poll_service(
145 db,
146 actor,
147 self.git_service.service_name,
148 lambda actor, since: self.git_service.fetch_recent_events(
149 actor=actor,
150 since=since,
151 db=db,
152 ),
153 )
154 return inserted
155
156 def _poll_service(self, db: Session, actor: str, service_name: str, fetcher) -> int:
157 state = db.scalar(
158 select(SyncState).where(SyncState.service == service_name).where(SyncState.actor == actor)
159 )
160 since = datetime.now(tz=UTC) - timedelta(days=30)
161 if state and state.cursor_value:
162 since = datetime.fromisoformat(state.cursor_value.replace("Z", "+00:00")).astimezone(UTC) - self._sync_overlap
163 logger.info(
164 "Using sync cursor for %s actor=%s with overlap; since=%s",
165 service_name,
166 actor,
167 since.isoformat(),
168 )
169
170 result = fetcher(actor=actor, since=since)
171 inserted = self._insert_events(db, result.events)
172 self._upsert_sync_state(db, service_name, actor, result.cursor)
173 logger.info("Polled %s for %s: inserted=%s", service_name, actor, inserted)
174 return inserted
175
176 @staticmethod
177 def _insert_events(db: Session, events: list[NormalizedEvent]) -> int:
178 inserted = 0
179 for event in events:
180 try:
181 with db.begin_nested():
182 model = ContributionEvent(**event.model_dump())
183 db.add(model)
184 db.flush()
185 inserted += 1
186 except IntegrityError:
187 logger.info("Skipping duplicate event %s for service %s", event.external_uid, event.service)
188 return inserted
189
190 @staticmethod
191 def _upsert_sync_state(db: Session, service: str, actor: str, cursor_value: str) -> None:
192 state = db.scalar(select(SyncState).where(SyncState.service == service).where(SyncState.actor == actor))
193 now = datetime.now(tz=UTC)
194 if state is None:
195 db.add(SyncState(service=service, actor=actor, cursor_value=cursor_value, updated_at=now))
196 db.flush()
197 return
198
199 state.cursor_value = cursor_value
200 state.updated_at = now
201 db.add(state)
202 db.flush()
203
204 def _sync_tracked_repositories(self, db: Session, actor: str) -> None:
205 configured = self.git_service.settings.git_tracked_repositories
206 for repo_name in configured:
207 canonical_repo_name = canonicalize_repository_name(actor, repo_name)
208 existing = db.scalar(
209 select(TrackedRepository)
210 .where(TrackedRepository.service == self.git_service.service_name)
211 .where(TrackedRepository.actor == actor)
212 .where(TrackedRepository.repo_name == canonical_repo_name)
213 )
214 if existing is None:
215 db.add(
216 TrackedRepository(
217 service=self.git_service.service_name,
218 repo_name=canonical_repo_name,
219 actor=actor,
220 )
221 )
222 db.flush()
223
224 def _update_tracked_actor_poll_state(self, db: Session, actor: str, status: str, error: str | None) -> None:
225 tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
226 tracked_actor.last_poll_status = status
227 tracked_actor.last_poll_error = error
228 now = datetime.now(tz=UTC)
229 if status == "indexed":
230 tracked_actor.discovery_state = "indexed"
231 tracked_actor.last_polled_at = now
232 tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds)
233 tracked_actor.priority_boosted_at = None
234 tracked_actor.poll_attempts = 0
235 elif status == "error":
236 tracked_actor.discovery_state = "error"
237 backoff_seconds = min(
238 self.settings.discovery_error_backoff_seconds * max(tracked_actor.poll_attempts, 1),
239 self.settings.discovery_error_backoff_max_seconds,
240 )
241 tracked_actor.next_poll_after = now + timedelta(seconds=backoff_seconds)
242 tracked_actor.priority_boosted_at = None
243 db.add(tracked_actor)
244 db.flush()
245
246 def poll_recent_backfill(self, db: Session, actor: str) -> int:
247 try:
248 inserted = self._run_backfill_batches(db, actor)
249 except Exception as exc:
250 db.rollback()
251 self._update_tracked_actor_poll_state(db, actor, status="error", error=str(exc))
252 db.commit()
253 raise
254
255 self._schedule_pending_backfill_or_repoll(db, actor)
256 db.commit()
257 return inserted
258
259 def _schedule_pending_backfill_or_repoll(self, db: Session, actor: str) -> None:
260 tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
261 now = datetime.now(tz=UTC)
262 tracked_actor.discovery_state = "indexed"
263 tracked_actor.last_poll_status = "indexed"
264 tracked_actor.last_poll_error = None
265 tracked_actor.priority_boosted_at = None
266 tracked_actor.poll_attempts = 0
267 if tracked_actor.recent_backfill_status == "completed":
268 tracked_actor.next_poll_after = now + timedelta(seconds=self.settings.indexed_actor_repoll_seconds)
269 else:
270 tracked_actor.next_poll_after = now
271 db.add(tracked_actor)
272 db.flush()
273
274 def _run_backfill_batches(self, db: Session, actor: str) -> int:
275 tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
276 if tracked_actor.recent_backfill_status == "completed":
277 return 0
278
279 now = datetime.now(tz=UTC)
280 recent_since = now - timedelta(days=RETENTION_DAYS)
281 db.add(tracked_actor)
282 db.flush()
283 total_inserted = 0
284
285 recent_services = [
286 (self.todo_service.service_name, self.todo_service.fetch_recent_backfill_batch),
287 (self.git_service.service_name, self.git_service.fetch_recent_backfill_batch),
288 ]
289 if tracked_actor.recent_backfill_status != "completed":
290 if tracked_actor.recent_backfill_started_at is None:
291 tracked_actor.recent_backfill_started_at = now
292 tracked_actor.recent_backfill_status = "in_progress"
293 tracked_actor.last_recent_backfill_error = None
294 total_inserted += self._run_backfill_scope(
295 db,
296 actor=actor,
297 scope="recent",
298 services=recent_services,
299 batches_per_service=RECENT_BACKFILL_BATCHES_PER_SERVICE,
300 since=recent_since,
301 )
302 recent_statuses = db.scalars(
303 select(ServiceBackfillState.status)
304 .where(ServiceBackfillState.actor == actor)
305 .where(ServiceBackfillState.scope == "recent")
306 ).all()
307 if recent_statuses and all(status == "completed" for status in recent_statuses):
308 tracked_actor.recent_backfill_status = "completed"
309 tracked_actor.recent_backfill_completed_at = datetime.now(tz=UTC)
310 tracked_actor.last_recent_backfill_error = None
311 db.add(tracked_actor)
312 db.flush()
313
314 return total_inserted
315
316 def _run_backfill_scope(
317 self,
318 db: Session,
319 *,
320 actor: str,
321 scope: str,
322 services,
323 batches_per_service: int,
324 since: datetime | None,
325 ) -> int:
326 tracked_actor = self.track_actor_request(db, actor, update_last_requested=False)
327 total_inserted = 0
328 for service_name, fetcher in services:
329 state = db.scalar(
330 select(ServiceBackfillState)
331 .where(ServiceBackfillState.actor == actor)
332 .where(ServiceBackfillState.service == service_name)
333 .where(ServiceBackfillState.scope == scope)
334 )
335 if state is None:
336 state = ServiceBackfillState(
337 actor=actor,
338 service=service_name,
339 scope=scope,
340 cursor_json=None,
341 status="pending",
342 started_at=None,
343 completed_at=None,
344 last_error=None,
345 updated_at=datetime.now(tz=UTC),
346 )
347 db.add(state)
348 db.flush()
349
350 if state.status == "completed":
351 continue
352
353 if state.started_at is None:
354 state.started_at = datetime.now(tz=UTC)
355 state.status = "in_progress"
356 state.updated_at = datetime.now(tz=UTC)
357
358 for _ in range(batches_per_service):
359 try:
360 if since is None:
361 result = fetcher(actor=actor, cursor_state=state.cursor_json)
362 else:
363 result = fetcher(actor=actor, cursor_state=state.cursor_json, since=since)
364 inserted = self._insert_events(db, result.events)
365 total_inserted += inserted
366 state.cursor_json = copy.deepcopy(result.cursor_state)
367 state.last_error = None
368 state.updated_at = datetime.now(tz=UTC)
369 if result.complete:
370 state.status = "completed"
371 state.completed_at = datetime.now(tz=UTC)
372 logger.info(
373 "Backfill complete for scope=%s service=%s actor=%s inserted=%s",
374 scope,
375 service_name,
376 actor,
377 inserted,
378 )
379 break
380 logger.info(
381 "Backfill batch complete for scope=%s service=%s actor=%s inserted=%s",
382 scope,
383 service_name,
384 actor,
385 inserted,
386 )
387 except Exception as exc:
388 state.status = "error"
389 state.last_error = str(exc)
390 state.updated_at = datetime.now(tz=UTC)
391 tracked_actor.recent_backfill_status = "error"
392 tracked_actor.last_recent_backfill_error = str(exc)
393 db.add(state)
394 db.add(tracked_actor)
395 db.flush()
396 raise
397
398 db.add(state)
399 db.flush()
400
401 return total_inserted
402
403 def prune_old_events(self, db: Session) -> int:
404 deleted = prune_contribution_events(db)
405 db.commit()
406 return deleted