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