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.py91
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