aboutsummaryrefslogtreecommitdiff
path: root/src/srht_contrib/jobs/poller.py
diff options
context:
space:
mode:
Diffstat (limited to 'src/srht_contrib/jobs/poller.py')
-rw-r--r--src/srht_contrib/jobs/poller.py173
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