summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/srht_contrib/jobs/poller.py91
-rw-r--r--src/srht_contrib/models.py19
-rw-r--r--src/srht_contrib/schemas.py3
-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
7 files changed, 352 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
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