From e1a3233ec2204055faa8dadb18cb2efa12cf0e17 Mon Sep 17 00:00:00 2001 From: Christian Cleberg Date: Sat, 11 Apr 2026 17:39:21 -0500 Subject: feat: add resumable historical backfill for actor activity --- src/srht_contrib/jobs/poller.py | 91 ++++++++++++++++++++- src/srht_contrib/models.py | 19 +++++ src/srht_contrib/schemas.py | 3 + src/srht_contrib/services/aggregator.py | 3 + src/srht_contrib/services/git.py | 90 +++++++++++++++++++++ src/srht_contrib/services/todo.py | 136 ++++++++++++++++++++++++++++++++ src/srht_contrib/services/types.py | 12 +++ 7 files changed, 352 insertions(+), 2 deletions(-) create mode 100644 src/srht_contrib/services/types.py (limited to 'src') 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 diff --git a/src/srht_contrib/models.py b/src/srht_contrib/models.py index 2f33482..a99a306 100644 --- a/src/srht_contrib/models.py +++ b/src/srht_contrib/models.py @@ -71,3 +71,22 @@ 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) + 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) + last_backfill_error: Mapped[str | None] = mapped_column(Text, nullable=True) + + +class ServiceBackfillState(Base): + __tablename__ = "service_backfill_states" + __table_args__ = (UniqueConstraint("actor", "service", name="uq_service_backfill_state_actor_service"),) + + 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) + 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) + completed_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + last_error: Mapped[str | None] = mapped_column(Text, nullable=True) + updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False) diff --git a/src/srht_contrib/schemas.py b/src/srht_contrib/schemas.py index bebf93d..6c77f6b 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_backfilled: bool = False + backfill_state: Literal["pending", "in_progress", "completed", "error"] = "pending" + backfill_completed_at: datetime | None = None class ContributionCalendarResponse(ContributionIndexMetadata): diff --git a/src/srht_contrib/services/aggregator.py b/src/srht_contrib/services/aggregator.py index 90ef93b..1c8af86 100644 --- a/src/srht_contrib/services/aggregator.py +++ b/src/srht_contrib/services/aggregator.py @@ -81,6 +81,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_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, ) def _query_daily_aggregates(self, db: Session, actor: str, start: date, end: date) -> list[DailyAggregate]: diff --git a/src/srht_contrib/services/git.py b/src/srht_contrib/services/git.py index 9531709..69a6c1c 100644 --- a/src/srht_contrib/services/git.py +++ b/src/srht_contrib/services/git.py @@ -8,6 +8,7 @@ from typing import Any from srht_contrib.config import Settings from srht_contrib.schemas import NormalizedEvent from srht_contrib.services.srht_client import SourceHutGraphQLClient +from srht_contrib.services.types import BackfillBatchResult from srht_contrib.utils.dates import ensure_utc, parse_datetime from srht_contrib.utils.identity import ActorIdentityResolver @@ -117,6 +118,95 @@ class GitIngestionService: ) return repositories + def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: + state = { + "discovery_cursor": None, + "discovery_complete": False, + "repository_queue": sorted( + { + self._canonical_repository_name(actor, repository) + for repository in self.settings.git_tracked_repositories + } + ), + "current_repository": None, + } + if cursor_state: + state.update(cursor_state) + + if not state["discovery_complete"]: + data = self.client.execute( + USER_REPOSITORIES_QUERY, + {"username": actor.lstrip("~"), "cursor": state["discovery_cursor"]}, + ) + user = data.get("user") or {} + repositories_page = user.get("repositories") or {} + results = repositories_page.get("results") or [] + state["discovery_cursor"] = repositories_page.get("cursor") + state["discovery_complete"] = not bool(state["discovery_cursor"]) + known = set(state["repository_queue"]) + current_repository = state.get("current_repository") + if current_repository: + known.add(current_repository["name"]) + for repository in results: + if not isinstance(repository, dict): + continue + name = repository.get("name") + repository_owner = ((repository.get("owner") or {}).get("canonicalName") or actor).strip() + if not name or not repository_owner: + continue + canonical_name = f"{repository_owner}/{name}" + if canonical_name not in known: + state["repository_queue"].append(canonical_name) + known.add(canonical_name) + state["repository_queue"] = sorted(state["repository_queue"]) + logger.info( + "git backfill discovery actor=%s page_count=%s queue=%s next_cursor=%s", + actor, + len(results), + len(state["repository_queue"]), + bool(state["discovery_cursor"]), + ) + complete = state["discovery_complete"] and not state["repository_queue"] and not state["current_repository"] + return BackfillBatchResult(events=[], cursor_state=state, complete=complete) + + if state["current_repository"] is None: + if not state["repository_queue"]: + return BackfillBatchResult(events=[], cursor_state=None, complete=True) + state["current_repository"] = {"name": state["repository_queue"].pop(0), "cursor": None} + + repository_name = state["current_repository"]["name"] + owner, repo_name = self._split_repository(actor, repository_name) + data = self.client.execute( + REPOSITORY_LOG_QUERY, + {"username": owner, "repoName": repo_name, "cursor": state["current_repository"]["cursor"]}, + ) + user = data.get("user") or {} + repository = user.get("repository") or {} + 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 + ] + logger.info( + "git backfill actor=%s repository=%s commits=%s next_cursor=%s", + actor, + repository_name, + len(commits), + bool(next_cursor), + ) + if next_cursor: + state["current_repository"]["cursor"] = next_cursor + else: + state["current_repository"] = None + + complete = state["discovery_complete"] and not state["repository_queue"] and not state["current_repository"] + return BackfillBatchResult(events=events, cursor_state=None if complete else state, complete=complete) + def _discover_owned_repositories(self, actor: str) -> list[str]: owner = actor.lstrip("~") repositories: list[str] = [] diff --git a/src/srht_contrib/services/todo.py b/src/srht_contrib/services/todo.py index 2d25acf..0142005 100644 --- a/src/srht_contrib/services/todo.py +++ b/src/srht_contrib/services/todo.py @@ -8,6 +8,7 @@ from typing import Any from srht_contrib.config import Settings from srht_contrib.schemas import NormalizedEvent from srht_contrib.services.srht_client import SourceHutGraphQLClient +from srht_contrib.services.types import BackfillBatchResult from srht_contrib.utils.dates import ensure_utc, parse_datetime @@ -293,6 +294,141 @@ class TodoIngestionService: tracker_events = self._fetch_from_trackers(actor=actor, since=since_dt) return TodoPollResult(events=tracker_events, cursor=cursor_time) + def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult: + state = { + "trackers_cursor": None, + "tracker_queue": [], + "current_tracker": None, + "current_ticket": None, + "trackers_loaded": False, + } + if cursor_state: + state.update(cursor_state) + + if not state["trackers_loaded"]: + data = self.client.execute(TODO_TRACKERS_QUERY, {"cursor": state["trackers_cursor"]}) + me = data.get("me") or {} + trackers_page = me.get("trackers") or {} + results = trackers_page.get("results") or [] + state["trackers_cursor"] = trackers_page.get("cursor") + for tracker in results: + if not isinstance(tracker, dict): + continue + state["tracker_queue"].append( + { + "id": str(tracker.get("id")), + "rid": str(tracker.get("rid")), + "name": tracker.get("name"), + "tickets_cursor": None, + "pending_tickets": [], + } + ) + state["trackers_loaded"] = not bool(state["trackers_cursor"]) + logger.info( + "todo backfill tracker discovery actor=%s trackers=%s next_cursor=%s queue=%s", + actor, + len(results), + bool(state["trackers_cursor"]), + len(state["tracker_queue"]), + ) + complete = state["trackers_loaded"] and not state["tracker_queue"] + return BackfillBatchResult(events=[], cursor_state=None if complete else state, complete=complete) + + if state["current_tracker"] is None: + if not state["tracker_queue"]: + return BackfillBatchResult(events=[], cursor_state=None, complete=True) + state["current_tracker"] = state["tracker_queue"].pop(0) + + current_tracker = state["current_tracker"] + if state["current_ticket"] is None and not current_tracker["pending_tickets"]: + data = self.client.execute( + TODO_TRACKER_TICKETS_QUERY, + {"trackerRid": current_tracker["rid"], "cursor": current_tracker["tickets_cursor"]}, + ) + 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)] + ) + 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"), + len(tickets), + bool(current_tracker["tickets_cursor"]), + len(current_tracker["pending_tickets"]), + ) + if not current_tracker["pending_tickets"] and not current_tracker["tickets_cursor"]: + state["current_tracker"] = None + return BackfillBatchResult(events=[], cursor_state=state, complete=False) + + if state["current_ticket"] is None: + if current_tracker["pending_tickets"]: + next_ticket = current_tracker["pending_tickets"].pop(0) + state["current_ticket"] = {**next_ticket, "cursor": None} + else: + state["current_tracker"] = None + return BackfillBatchResult(events=[], cursor_state=state, complete=False) + + current_ticket = state["current_ticket"] + data = self.client.execute( + TODO_TICKET_EVENTS_QUERY, + {"trackerRid": current_tracker["rid"], "ticketId": current_ticket["id"], "cursor": current_ticket["cursor"]}, + ) + tracker = data.get("tracker") or {} + ticket_payload = (tracker.get("ticket") or {}) if isinstance(tracker, dict) else {} + event_page = ticket_payload.get("events") or {} + page_events = event_page.get("results") or [] + next_cursor = event_page.get("cursor") + logger.info( + "todo backfill ticket events ref=%s tracker=%s count=%s next_cursor=%s", + current_ticket["ref"], + current_tracker.get("name") or current_tracker.get("id"), + len(page_events), + bool(next_cursor), + ) + + events: list[NormalizedEvent] = [] + for event in page_events: + if not isinstance(event, dict): + continue + event["ticket"] = { + "id": ticket_payload.get("id", current_ticket["id"]), + "ref": ticket_payload.get("ref", current_ticket["ref"]), + "status": ticket_payload.get("status"), + "resolution": ticket_payload.get("resolution"), + "tracker": {"name": current_tracker.get("name")}, + } + occurred_at = parse_datetime(event["created"]) + for change in event.get("changes") or []: + if not isinstance(change, dict): + continue + normalized = _normalize_event_change( + settings=self.settings, + actor=actor, + event=event, + change=change, + occurred_at=occurred_at, + ) + if normalized is not None: + events.append(normalized) + + if next_cursor: + state["current_ticket"]["cursor"] = next_cursor + else: + state["current_ticket"] = None + if not current_tracker["pending_tickets"] and not current_tracker["tickets_cursor"]: + state["current_tracker"] = None + + complete = ( + state["trackers_loaded"] + and not state["tracker_queue"] + and state["current_tracker"] is None + and state["current_ticket"] is None + ) + return BackfillBatchResult(events=events, cursor_state=None if complete else state, complete=complete) + def _fetch_from_activity_feed(self, actor: str, since: datetime) -> TodoPollResult: events: list[NormalizedEvent] = [] cursor: str | None = None diff --git a/src/srht_contrib/services/types.py b/src/srht_contrib/services/types.py new file mode 100644 index 0000000..69e64af --- /dev/null +++ b/src/srht_contrib/services/types.py @@ -0,0 +1,12 @@ +from __future__ import annotations + +from dataclasses import dataclass + +from srht_contrib.schemas import NormalizedEvent + + +@dataclass(slots=True) +class BackfillBatchResult: + events: list[NormalizedEvent] + cursor_state: dict | None + complete: bool -- cgit v1.2.3