diff options
| author | Christian Cleberg <[email protected]> | 2026-04-11 17:39:21 -0500 |
|---|---|---|
| committer | Christian Cleberg <[email protected]> | 2026-04-11 17:39:21 -0500 |
| commit | e1a3233ec2204055faa8dadb18cb2efa12cf0e17 (patch) | |
| tree | fa6b70671c3f6b71a690e1c99f368616cb1a8e4d /src/srht_contrib/jobs/poller.py | |
| parent | aebe493dcd2b46bf3ae46543f2c0e81625d93ae5 (diff) | |
| download | hutch-stats-e1a3233ec2204055faa8dadb18cb2efa12cf0e17.tar.gz hutch-stats-e1a3233ec2204055faa8dadb18cb2efa12cf0e17.tar.bz2 hutch-stats-e1a3233ec2204055faa8dadb18cb2efa12cf0e17.zip | |
feat: add resumable historical backfill for actor activity
Diffstat (limited to 'src/srht_contrib/jobs/poller.py')
| -rw-r--r-- | src/srht_contrib/jobs/poller.py | 91 |
1 files changed, 89 insertions, 2 deletions
diff --git a/src/srht_contrib/jobs/poller.py b/src/srht_contrib/jobs/poller.py index 910f1ab..2b00090 100644 --- a/src/srht_contrib/jobs/poller.py +++ b/src/srht_contrib/jobs/poller.py @@ -7,7 +7,7 @@ from sqlalchemy import select from sqlalchemy.exc import IntegrityError from sqlalchemy.orm import Session -from srht_contrib.models import ContributionEvent, SyncState, TrackedActor, TrackedRepository +from srht_contrib.models import ContributionEvent, ServiceBackfillState, SyncState, TrackedActor, TrackedRepository from srht_contrib.schemas import NormalizedEvent from srht_contrib.services.git import GitIngestionService from srht_contrib.services.srht_client import SourceHutClientError @@ -35,6 +35,7 @@ class PollerService: raise self._update_tracked_actor_poll_state(db, actor, status="indexed", error=None) + inserted += self._run_backfill_batches(db, actor) db.commit() return inserted @@ -61,7 +62,7 @@ class PollerService: def track_actor_request(self, db: Session, actor: str, *, update_last_requested: bool = True) -> TrackedActor: tracked_actor = db.scalar(select(TrackedActor).where(TrackedActor.actor == actor)) if tracked_actor is None: - tracked_actor = TrackedActor(actor=actor, is_active=True) + tracked_actor = TrackedActor(actor=actor, is_active=True, backfill_status="pending") db.add(tracked_actor) tracked_actor.is_active = True @@ -172,3 +173,89 @@ class PollerService: tracked_actor.last_polled_at = datetime.now(tz=UTC) db.add(tracked_actor) db.flush() + + def _run_backfill_batches(self, db: Session, actor: str) -> int: + tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) + if tracked_actor.backfill_status == "completed": + return 0 + + if tracked_actor.backfill_started_at is None: + tracked_actor.backfill_started_at = datetime.now(tz=UTC) + tracked_actor.backfill_status = "in_progress" + tracked_actor.last_backfill_error = None + db.add(tracked_actor) + db.flush() + total_inserted = 0 + + services = [ + (self.todo_service.service_name, self.todo_service.fetch_backfill_batch), + (self.git_service.service_name, self.git_service.fetch_backfill_batch), + ] + all_complete = True + for service_name, fetcher in services: + state = db.scalar( + select(ServiceBackfillState) + .where(ServiceBackfillState.actor == actor) + .where(ServiceBackfillState.service == service_name) + ) + if state is None: + state = ServiceBackfillState( + actor=actor, + service=service_name, + cursor_json=None, + status="pending", + started_at=None, + completed_at=None, + last_error=None, + updated_at=datetime.now(tz=UTC), + ) + db.add(state) + db.flush() + + if state.status == "completed": + continue + + all_complete = False + if state.started_at is None: + state.started_at = datetime.now(tz=UTC) + state.status = "in_progress" + state.updated_at = datetime.now(tz=UTC) + try: + result = fetcher(actor=actor, cursor_state=state.cursor_json) + inserted = self._insert_events(db, result.events) + total_inserted += inserted + state.cursor_json = result.cursor_state + state.last_error = None + state.updated_at = datetime.now(tz=UTC) + if result.complete: + state.status = "completed" + state.completed_at = datetime.now(tz=UTC) + logger.info("Backfill complete for service=%s actor=%s inserted=%s", service_name, actor, inserted) + else: + logger.info("Backfill batch complete for service=%s actor=%s inserted=%s", service_name, actor, inserted) + except Exception as exc: + state.status = "error" + state.last_error = str(exc) + state.updated_at = datetime.now(tz=UTC) + tracked_actor.backfill_status = "error" + tracked_actor.last_backfill_error = str(exc) + db.add(state) + db.add(tracked_actor) + db.flush() + raise + + db.add(state) + db.flush() + + completed = db.scalars( + select(ServiceBackfillState.status).where(ServiceBackfillState.actor == actor) + ).all() + if completed and all(status == "completed" for status in completed): + tracked_actor.backfill_status = "completed" + tracked_actor.backfill_completed_at = datetime.now(tz=UTC) + tracked_actor.last_backfill_error = None + elif tracked_actor.backfill_status != "error": + tracked_actor.backfill_status = "in_progress" + db.add(tracked_actor) + db.flush() + return total_inserted |
