diff options
| author | Christian Cleberg <[email protected]> | 2026-04-11 17:39:21 -0500 |
|---|---|---|
| committer | Christian Cleberg <[email protected]> | 2026-04-11 17:39:21 -0500 |
| commit | e1a3233ec2204055faa8dadb18cb2efa12cf0e17 (patch) | |
| tree | fa6b70671c3f6b71a690e1c99f368616cb1a8e4d /src/srht_contrib/services | |
| parent | aebe493dcd2b46bf3ae46543f2c0e81625d93ae5 (diff) | |
| download | hutch-stats-e1a3233ec2204055faa8dadb18cb2efa12cf0e17.tar.gz hutch-stats-e1a3233ec2204055faa8dadb18cb2efa12cf0e17.tar.bz2 hutch-stats-e1a3233ec2204055faa8dadb18cb2efa12cf0e17.zip | |
feat: add resumable historical backfill for actor activity
Diffstat (limited to 'src/srht_contrib/services')
| -rw-r--r-- | src/srht_contrib/services/aggregator.py | 3 | ||||
| -rw-r--r-- | src/srht_contrib/services/git.py | 90 | ||||
| -rw-r--r-- | src/srht_contrib/services/todo.py | 136 | ||||
| -rw-r--r-- | src/srht_contrib/services/types.py | 12 |
4 files changed, 241 insertions, 0 deletions
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 |
