diff options
Diffstat (limited to 'src/srht_contrib/jobs')
| -rw-r--r-- | src/srht_contrib/jobs/poller.py | 173 |
1 files changed, 128 insertions, 45 deletions
diff --git a/src/srht_contrib/jobs/poller.py b/src/srht_contrib/jobs/poller.py index 45992e4..c41f3d4 100644 --- a/src/srht_contrib/jobs/poller.py +++ b/src/srht_contrib/jobs/poller.py @@ -18,6 +18,9 @@ from srht_contrib.utils.repositories import canonicalize_repository_name logger = logging.getLogger(__name__) SYNC_OVERLAP = timedelta(hours=24) +RECENT_BACKFILL_DAYS = 365 +RECENT_BACKFILL_BATCHES_PER_SERVICE = 5 +FULL_BACKFILL_BATCHES_PER_SERVICE = 1 class PollerService: @@ -63,7 +66,12 @@ 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, backfill_status="pending") + tracked_actor = TrackedActor( + actor=actor, + is_active=True, + recent_backfill_status="pending", + backfill_status="pending", + ) db.add(tracked_actor) tracked_actor.is_active = True @@ -177,32 +185,98 @@ class PollerService: 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": + if tracked_actor.recent_backfill_status == "completed" and 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 + now = datetime.now(tz=UTC) + recent_since = now - timedelta(days=RECENT_BACKFILL_DAYS) 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), + recent_services = [ + (self.todo_service.service_name, self.todo_service.fetch_recent_backfill_batch), + (self.git_service.service_name, self.git_service.fetch_recent_backfill_batch), ] - all_complete = True + if tracked_actor.recent_backfill_status != "completed": + if tracked_actor.recent_backfill_started_at is None: + tracked_actor.recent_backfill_started_at = now + tracked_actor.recent_backfill_status = "in_progress" + tracked_actor.last_recent_backfill_error = None + total_inserted += self._run_backfill_scope( + db, + actor=actor, + scope="recent", + services=recent_services, + batches_per_service=RECENT_BACKFILL_BATCHES_PER_SERVICE, + since=recent_since, + ) + recent_statuses = db.scalars( + select(ServiceBackfillState.status) + .where(ServiceBackfillState.actor == actor) + .where(ServiceBackfillState.scope == "recent") + ).all() + if recent_statuses and all(status == "completed" for status in recent_statuses): + tracked_actor.recent_backfill_status = "completed" + tracked_actor.recent_backfill_completed_at = datetime.now(tz=UTC) + tracked_actor.last_recent_backfill_error = None + db.add(tracked_actor) + db.flush() + + if tracked_actor.recent_backfill_status == "completed" and tracked_actor.backfill_status != "completed": + if tracked_actor.backfill_started_at is None: + tracked_actor.backfill_started_at = now + tracked_actor.backfill_status = "in_progress" + tracked_actor.last_backfill_error = None + total_inserted += self._run_backfill_scope( + db, + actor=actor, + scope="full", + services=[ + (self.todo_service.service_name, self.todo_service.fetch_backfill_batch), + (self.git_service.service_name, self.git_service.fetch_backfill_batch), + ], + batches_per_service=FULL_BACKFILL_BATCHES_PER_SERVICE, + since=None, + ) + full_statuses = db.scalars( + select(ServiceBackfillState.status) + .where(ServiceBackfillState.actor == actor) + .where(ServiceBackfillState.scope == "full") + ).all() + if full_statuses and all(status == "completed" for status in full_statuses): + tracked_actor.backfill_status = "completed" + tracked_actor.backfill_completed_at = datetime.now(tz=UTC) + tracked_actor.last_backfill_error = None + db.add(tracked_actor) + db.flush() + + return total_inserted + + def _run_backfill_scope( + self, + db: Session, + *, + actor: str, + scope: str, + services, + batches_per_service: int, + since: datetime | None, + ) -> int: + tracked_actor = self.track_actor_request(db, actor, update_last_requested=False) + total_inserted = 0 for service_name, fetcher in services: state = db.scalar( select(ServiceBackfillState) .where(ServiceBackfillState.actor == actor) .where(ServiceBackfillState.service == service_name) + .where(ServiceBackfillState.scope == scope) ) if state is None: state = ServiceBackfillState( actor=actor, service=service_name, + scope=scope, cursor_json=None, status="pending", started_at=None, @@ -216,47 +290,56 @@ class PollerService: 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 = copy.deepcopy(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 + + for _ in range(batches_per_service): + try: + if since is None: + result = fetcher(actor=actor, cursor_state=state.cursor_json) + else: + result = fetcher(actor=actor, cursor_state=state.cursor_json, since=since) + inserted = self._insert_events(db, result.events) + total_inserted += inserted + state.cursor_json = copy.deepcopy(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 scope=%s service=%s actor=%s inserted=%s", + scope, + service_name, + actor, + inserted, + ) + break + logger.info( + "Backfill batch complete for scope=%s service=%s actor=%s inserted=%s", + scope, + service_name, + actor, + inserted, + ) + except Exception as exc: + state.status = "error" + state.last_error = str(exc) + state.updated_at = datetime.now(tz=UTC) + if scope == "recent": + tracked_actor.recent_backfill_status = "error" + tracked_actor.last_recent_backfill_error = str(exc) + else: + 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 |
