summaryrefslogtreecommitdiff
path: root/src/srht_contrib/services
diff options
context:
space:
mode:
Diffstat (limited to 'src/srht_contrib/services')
-rw-r--r--src/srht_contrib/services/aggregator.py3
-rw-r--r--src/srht_contrib/services/git.py90
-rw-r--r--src/srht_contrib/services/todo.py136
-rw-r--r--src/srht_contrib/services/types.py12
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