diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/srht_contrib/jobs/poller.py | 173 | ||||
| -rw-r--r-- | src/srht_contrib/models.py | 7 | ||||
| -rw-r--r-- | src/srht_contrib/schemas.py | 3 | ||||
| -rw-r--r-- | src/srht_contrib/services/aggregator.py | 9 | ||||
| -rw-r--r-- | src/srht_contrib/services/git.py | 39 | ||||
| -rw-r--r-- | src/srht_contrib/services/todo.py | 42 |
6 files changed, 214 insertions, 59 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 diff --git a/src/srht_contrib/models.py b/src/srht_contrib/models.py index a99a306..475bd22 100644 --- a/src/srht_contrib/models.py +++ b/src/srht_contrib/models.py @@ -71,6 +71,10 @@ class TrackedActor(Base): last_polled_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) last_poll_status: Mapped[str | None] = mapped_column(String(32), nullable=True) last_poll_error: Mapped[str | None] = mapped_column(Text, nullable=True) + recent_backfill_status: Mapped[str] = mapped_column(String(32), nullable=False, default="pending") + recent_backfill_started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + recent_backfill_completed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + last_recent_backfill_error: Mapped[str | None] = mapped_column(Text, nullable=True) backfill_status: Mapped[str] = mapped_column(String(32), nullable=False, default="pending") backfill_started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) backfill_completed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) @@ -79,11 +83,12 @@ class TrackedActor(Base): class ServiceBackfillState(Base): __tablename__ = "service_backfill_states" - __table_args__ = (UniqueConstraint("actor", "service", name="uq_service_backfill_state_actor_service"),) + __table_args__ = (UniqueConstraint("actor", "service", "scope", name="uq_service_backfill_state_actor_service_scope"),) id: Mapped[int] = mapped_column(Integer, primary_key=True) actor: Mapped[str] = mapped_column(String(255), nullable=False) service: Mapped[str] = mapped_column(String(32), nullable=False) + scope: Mapped[str] = mapped_column(String(16), nullable=False, default="full") cursor_json: Mapped[dict | None] = mapped_column(JSON, nullable=True) status: Mapped[str] = mapped_column(String(32), nullable=False, default="pending") started_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) diff --git a/src/srht_contrib/schemas.py b/src/srht_contrib/schemas.py index 6c77f6b..98d2a74 100644 --- a/src/srht_contrib/schemas.py +++ b/src/srht_contrib/schemas.py @@ -28,6 +28,9 @@ class ContributionIndexMetadata(BaseModel): is_indexed: bool last_polled_at: datetime | None = None indexing_state: Literal["pending", "indexed", "error"] + is_recent_window_backfilled: bool = False + recent_backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending" + recent_backfill_completed_at: datetime | None = None is_backfilled: bool = False backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending" backfill_completed_at: datetime | None = None diff --git a/src/srht_contrib/services/aggregator.py b/src/srht_contrib/services/aggregator.py index 1c8af86..4973f16 100644 --- a/src/srht_contrib/services/aggregator.py +++ b/src/srht_contrib/services/aggregator.py @@ -62,6 +62,12 @@ class ContributionAggregator: is_indexed=calendar.is_indexed, last_polled_at=calendar.last_polled_at, indexing_state=calendar.indexing_state, + is_recent_window_backfilled=calendar.is_recent_window_backfilled, + recent_backfill_state=calendar.recent_backfill_state, + recent_backfill_completed_at=calendar.recent_backfill_completed_at, + is_backfilled=calendar.is_backfilled, + backfill_state=calendar.backfill_state, + backfill_completed_at=calendar.backfill_completed_at, ) def _index_metadata(self, db: Session, actor: str) -> ContributionIndexMetadata: @@ -81,6 +87,9 @@ class ContributionAggregator: is_indexed=is_indexed, last_polled_at=tracked_actor.last_polled_at if tracked_actor is not None else None, indexing_state=indexing_state, + is_recent_window_backfilled=(tracked_actor.recent_backfill_status == "completed") if tracked_actor is not None else False, + recent_backfill_state=(tracked_actor.recent_backfill_status if tracked_actor is not None else "pending"), + recent_backfill_completed_at=tracked_actor.recent_backfill_completed_at if tracked_actor is not None else None, is_backfilled=(tracked_actor.backfill_status == "completed") if tracked_actor is not None else False, backfill_state=(tracked_actor.backfill_status if tracked_actor is not None else "pending"), backfill_completed_at=tracked_actor.backfill_completed_at if tracked_actor is not None else None, diff --git a/src/srht_contrib/services/git.py b/src/srht_contrib/services/git.py index ebb823e..bf7493b 100644 --- a/src/srht_contrib/services/git.py +++ b/src/srht_contrib/services/git.py @@ -120,6 +120,24 @@ class GitIngestionService: return repositories def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: + return self._fetch_backfill_batch(actor, cursor_state, since=None) + + def fetch_recent_backfill_batch( + self, + actor: str, + cursor_state: dict | None = None, + *, + since: datetime, + ) -> BackfillBatchResult: + return self._fetch_backfill_batch(actor, cursor_state, since=since) + + def _fetch_backfill_batch( + self, + actor: str, + cursor_state: dict | None, + *, + since: datetime | None, + ) -> BackfillBatchResult: state = { "discovery_cursor": None, "discovery_complete": False, @@ -186,13 +204,18 @@ class GitIngestionService: log_page = repository.get("log") or {} commits = log_page.get("results") or [] next_cursor = log_page.get("cursor") - events = [ - normalized - for commit in commits - if isinstance(commit, dict) - for normalized in [self._normalize_commit(actor=actor, repo_name=repo_name, commit=commit)] - if normalized is not None - ] + events: list[NormalizedEvent] = [] + stop_repository = False + for commit in commits: + if not isinstance(commit, dict): + continue + commit_time = parse_datetime((commit.get("author") or {}).get("time")) + if since is not None and commit_time < since: + stop_repository = True + break + normalized = self._normalize_commit(actor=actor, repo_name=repo_name, commit=commit) + if normalized is not None: + events.append(normalized) logger.info( "git backfill actor=%s repository=%s commits=%s next_cursor=%s", actor, @@ -200,7 +223,7 @@ class GitIngestionService: len(commits), bool(next_cursor), ) - if next_cursor: + if next_cursor and not stop_repository: state["current_repository"]["cursor"] = next_cursor else: state["current_repository"] = None diff --git a/src/srht_contrib/services/todo.py b/src/srht_contrib/services/todo.py index b7209b1..9f2ece6 100644 --- a/src/srht_contrib/services/todo.py +++ b/src/srht_contrib/services/todo.py @@ -296,6 +296,24 @@ class TodoIngestionService: return TodoPollResult(events=tracker_events, cursor=cursor_time) def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: + return self._fetch_backfill_batch(actor, cursor_state, since=None) + + def fetch_recent_backfill_batch( + self, + actor: str, + cursor_state: dict | None = None, + *, + since: datetime, + ) -> BackfillBatchResult: + return self._fetch_backfill_batch(actor, cursor_state, since=since) + + def _fetch_backfill_batch( + self, + actor: str, + cursor_state: dict | None, + *, + since: datetime | None, + ) -> BackfillBatchResult: state = { "trackers_cursor": None, "tracker_queue": [], @@ -349,10 +367,20 @@ class TodoIngestionService: tracker = data.get("tracker") or {} tickets_page = tracker.get("tickets") or {} tickets = tickets_page.get("results") or [] - current_tracker["tickets_cursor"] = tickets_page.get("cursor") - current_tracker["pending_tickets"].extend( - [{"id": int(ticket["id"]), "ref": str(ticket.get("ref") or ticket["id"])} for ticket in tickets if isinstance(ticket, dict)] - ) + next_tickets_cursor = tickets_page.get("cursor") + stop_tracker_paging = False + for ticket in tickets: + if not isinstance(ticket, dict): + continue + if since is not None: + updated = parse_datetime(ticket["updated"]) + if updated < since: + stop_tracker_paging = True + continue + current_tracker["pending_tickets"].append( + {"id": int(ticket["id"]), "ref": str(ticket.get("ref") or ticket["id"])} + ) + current_tracker["tickets_cursor"] = None if stop_tracker_paging else next_tickets_cursor logger.info( "todo backfill tracker=%s ticket page count=%s next_cursor=%s pending_tickets=%s", current_tracker.get("name") or current_tracker.get("id"), @@ -391,6 +419,7 @@ class TodoIngestionService: ) events: list[NormalizedEvent] = [] + stop_ticket_paging = False for event in page_events: if not isinstance(event, dict): continue @@ -402,6 +431,9 @@ class TodoIngestionService: "tracker": {"name": current_tracker.get("name")}, } occurred_at = parse_datetime(event["created"]) + if since is not None and occurred_at < since: + stop_ticket_paging = True + continue for change in event.get("changes") or []: if not isinstance(change, dict): continue @@ -415,7 +447,7 @@ class TodoIngestionService: if normalized is not None: events.append(normalized) - if next_cursor: + if next_cursor and not stop_ticket_paging: state["current_ticket"]["cursor"] = next_cursor else: state["current_ticket"] = None |
