summaryrefslogtreecommitdiff
path: root/src/srht_contrib/services
diff options
context:
space:
mode:
authorChristian Cleberg <[email protected]>2026-04-09 19:45:45 -0500
committerChristian Cleberg <[email protected]>2026-04-09 19:45:45 -0500
commitacbff854f2da96bddcaede1385e7fefeba0fb34b (patch)
tree48d4707cd3276d2825370f0793008a6ffcfff936 /src/srht_contrib/services
downloadhutch-stats-acbff854f2da96bddcaede1385e7fefeba0fb34b.tar.gz
hutch-stats-acbff854f2da96bddcaede1385e7fefeba0fb34b.tar.bz2
hutch-stats-acbff854f2da96bddcaede1385e7fefeba0fb34b.zip
initial commit
Diffstat (limited to 'src/srht_contrib/services')
-rw-r--r--src/srht_contrib/services/__init__.py1
-rw-r--r--src/srht_contrib/services/aggregator.py100
-rw-r--r--src/srht_contrib/services/git.py202
-rw-r--r--src/srht_contrib/services/srht_client.py77
-rw-r--r--src/srht_contrib/services/todo.py559
5 files changed, 939 insertions, 0 deletions
diff --git a/src/srht_contrib/services/__init__.py b/src/srht_contrib/services/__init__.py
new file mode 100644
index 0000000..02dea84
--- /dev/null
+++ b/src/srht_contrib/services/__init__.py
@@ -0,0 +1 @@
+"""Service layer."""
diff --git a/src/srht_contrib/services/aggregator.py b/src/srht_contrib/services/aggregator.py
new file mode 100644
index 0000000..800329a
--- /dev/null
+++ b/src/srht_contrib/services/aggregator.py
@@ -0,0 +1,100 @@
+from __future__ import annotations
+
+from dataclasses import dataclass
+from datetime import date
+
+from sqlalchemy import func, select
+from sqlalchemy.orm import Session
+
+from srht_contrib.models import ContributionEvent
+from srht_contrib.schemas import ContributionCalendarResponse, ContributionDay, ContributionStatsResponse
+from srht_contrib.utils.dates import date_range, date_to_utc_bounds
+
+
+@dataclass(slots=True)
+class DailyAggregate:
+ date: date
+ count: int
+ score: float
+
+
+class ContributionAggregator:
+ def build_calendar(self, db: Session, actor: str, start: date, end: date) -> ContributionCalendarResponse:
+ aggregates = self._query_daily_aggregates(db, actor, start, end)
+ by_day = {row.date: row for row in aggregates}
+ days = [
+ ContributionDay(
+ date=day,
+ count=by_day.get(day, DailyAggregate(date=day, count=0, score=0.0)).count,
+ score=by_day.get(day, DailyAggregate(date=day, count=0, score=0.0)).score,
+ )
+ for day in date_range(start, end)
+ ]
+ return ContributionCalendarResponse(actor=actor, from_date=start, to_date=end, days=days)
+
+ def build_stats(self, db: Session, actor: str, start: date, end: date) -> ContributionStatsResponse:
+ calendar = self.build_calendar(db, actor, start, end)
+ active_days = [day for day in calendar.days if day.count > 0]
+ streaks = self._streak_lengths(calendar.days)
+ current_streak = self._current_streak(calendar.days)
+
+ return ContributionStatsResponse(
+ actor=actor,
+ from_date=start,
+ to_date=end,
+ total_events=sum(day.count for day in calendar.days),
+ total_score=round(sum(day.score for day in calendar.days), 2),
+ active_days=len(active_days),
+ longest_streak=max(streaks, default=0),
+ current_streak=current_streak,
+ )
+
+ def _query_daily_aggregates(self, db: Session, actor: str, start: date, end: date) -> list[DailyAggregate]:
+ start_dt, _ = date_to_utc_bounds(start)
+ _, end_dt = date_to_utc_bounds(end)
+
+ stmt = (
+ select(
+ func.date(ContributionEvent.occurred_at).label("day"),
+ func.count(ContributionEvent.id).label("count"),
+ func.coalesce(func.sum(ContributionEvent.weight), 0.0).label("score"),
+ )
+ .where(ContributionEvent.actor == actor)
+ .where(ContributionEvent.occurred_at >= start_dt)
+ .where(ContributionEvent.occurred_at <= end_dt)
+ .group_by(func.date(ContributionEvent.occurred_at))
+ .order_by(func.date(ContributionEvent.occurred_at))
+ )
+ rows = db.execute(stmt).all()
+ return [
+ DailyAggregate(
+ date=date.fromisoformat(str(row.day)),
+ count=int(row.count),
+ score=round(float(row.score), 2),
+ )
+ for row in rows
+ ]
+
+ @staticmethod
+ def _streak_lengths(days: list[ContributionDay]) -> list[int]:
+ streaks: list[int] = []
+ current = 0
+ for day in days:
+ if day.count > 0:
+ current += 1
+ elif current > 0:
+ streaks.append(current)
+ current = 0
+ if current > 0:
+ streaks.append(current)
+ return streaks
+
+ @staticmethod
+ def _current_streak(days: list[ContributionDay]) -> int:
+ streak = 0
+ for day in reversed(days):
+ if day.count > 0:
+ streak += 1
+ else:
+ break
+ return streak
diff --git a/src/srht_contrib/services/git.py b/src/srht_contrib/services/git.py
new file mode 100644
index 0000000..20a18fe
--- /dev/null
+++ b/src/srht_contrib/services/git.py
@@ -0,0 +1,202 @@
+from __future__ import annotations
+
+from dataclasses import dataclass
+from datetime import UTC, datetime, timedelta
+import logging
+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.utils.dates import ensure_utc, parse_datetime
+from srht_contrib.utils.identity import ActorIdentityResolver
+
+
+logger = logging.getLogger(__name__)
+
+
+REPOSITORY_LOG_QUERY = """
+query RepositoryLog($username: String!, $repoName: String!, $cursor: Cursor) {
+ user(username: $username) {
+ repository(name: $repoName) {
+ name
+ owner {
+ canonicalName
+ }
+ log(cursor: $cursor) {
+ results {
+ id
+ shortId
+ author {
+ name
+ email
+ time
+ }
+ committer {
+ name
+ email
+ time
+ }
+ message
+ }
+ cursor
+ }
+ }
+ }
+}
+""".strip()
+
+
+@dataclass(slots=True)
+class GitPollResult:
+ events: list[NormalizedEvent]
+ cursor: str
+
+
+class GitIngestionService:
+ """Polls tracked git.sr.ht repositories and normalizes commits for one actor."""
+
+ service_name = "git"
+
+ def __init__(self, client: SourceHutGraphQLClient, settings: Settings) -> None:
+ self.client = client
+ self.settings = settings
+ self.identity_resolver = ActorIdentityResolver(settings.actor_aliases_json)
+
+ def fetch_recent_events(
+ self,
+ actor: str,
+ since: datetime | None = None,
+ repositories: list[str] | None = None,
+ ) -> GitPollResult:
+ since_dt = ensure_utc(since or (datetime.now(tz=UTC) - timedelta(days=30)))
+ tracked_repositories = repositories or self._tracked_repositories(actor)
+ if not tracked_repositories:
+ logger.info("git poll skipped for actor=%s because no tracked repositories are configured", actor)
+ return GitPollResult(events=[], cursor=datetime.now(tz=UTC).isoformat())
+
+ events: list[NormalizedEvent] = []
+ for repository in tracked_repositories:
+ owner, repo_name = self._split_repository(actor, repository)
+ repo_events = self._fetch_repository_commits(actor=actor, owner=owner, repo_name=repo_name, since=since_dt)
+ events.extend(repo_events)
+
+ logger.info("git poll complete for actor=%s normalized_events=%s", actor, len(events))
+ return GitPollResult(events=events, cursor=datetime.now(tz=UTC).isoformat())
+
+ def _tracked_repositories(self, actor: str) -> list[str]:
+ repositories = self.settings.git_tracked_repositories
+ return repositories
+
+ @staticmethod
+ def _split_repository(default_actor: str, repository: str) -> tuple[str, str]:
+ if "/" in repository:
+ owner, repo_name = repository.split("/", 1)
+ canonical_owner = owner if owner.startswith("~") else f"~{owner}"
+ return canonical_owner.lstrip("~"), repo_name
+ return default_actor.lstrip("~"), repository
+
+ def _fetch_repository_commits(
+ self,
+ *,
+ actor: str,
+ owner: str,
+ repo_name: str,
+ since: datetime,
+ ) -> list[NormalizedEvent]:
+ events: list[NormalizedEvent] = []
+ cursor: str | None = None
+
+ for _ in range(50):
+ data = self.client.execute(
+ REPOSITORY_LOG_QUERY,
+ {"username": owner, "repoName": repo_name, "cursor": cursor},
+ )
+ user = data.get("user") or {}
+ repository = user.get("repository") or {}
+ log_page = repository.get("log") or {}
+ commits = log_page.get("results") or []
+ cursor = log_page.get("cursor")
+ logger.info(
+ "git repository=%s/%s commit page count=%s next_cursor=%s",
+ owner,
+ repo_name,
+ len(commits),
+ bool(cursor),
+ )
+
+ stop_paging = False
+ for commit in commits:
+ if not isinstance(commit, dict):
+ continue
+ commit_time = parse_datetime((commit.get("author") or {}).get("time"))
+ if commit_time < since:
+ stop_paging = True
+ logger.info(
+ "git commit %s skipped because commit_time=%s is before since=%s",
+ commit.get("shortId") or commit.get("id"),
+ commit_time.isoformat(),
+ since.isoformat(),
+ )
+ continue
+
+ normalized = self._normalize_commit(actor=actor, repo_name=repo_name, commit=commit)
+ if normalized is not None:
+ logger.info(
+ "git commit accepted repo=%s shortId=%s author=%s email=%s",
+ repo_name,
+ commit.get("shortId"),
+ (commit.get("author") or {}).get("name"),
+ (commit.get("author") or {}).get("email"),
+ )
+ events.append(normalized)
+ else:
+ logger.info(
+ "git commit skipped repo=%s shortId=%s author=%s email=%s",
+ repo_name,
+ commit.get("shortId"),
+ (commit.get("author") or {}).get("name"),
+ (commit.get("author") or {}).get("email"),
+ )
+
+ if stop_paging or not cursor:
+ break
+
+ return events
+
+ def _normalize_commit(
+ self,
+ *,
+ actor: str,
+ repo_name: str,
+ commit: dict[str, Any],
+ ) -> NormalizedEvent | None:
+ author = commit.get("author") or {}
+ candidate_aliases = [
+ actor,
+ author.get("email", ""),
+ author.get("name", ""),
+ ]
+ matched_actor = None
+ for candidate in candidate_aliases:
+ canonical = self.identity_resolver.canonicalize(candidate)
+ if canonical == actor:
+ matched_actor = canonical
+ break
+
+ if matched_actor is None:
+ return None
+
+ commit_id = str(commit["id"])
+ commit_time = parse_datetime(author["time"])
+ return NormalizedEvent(
+ service=self.service_name,
+ event_type="commit",
+ actor=matched_actor,
+ repo_name=repo_name,
+ resource_id=commit_id,
+ external_uid=f"git:{repo_name}:{commit_id}",
+ occurred_at=commit_time,
+ weight=self.settings.event_weights["commit"],
+ raw_payload_json=commit,
+ )
diff --git a/src/srht_contrib/services/srht_client.py b/src/srht_contrib/services/srht_client.py
new file mode 100644
index 0000000..c8e2649
--- /dev/null
+++ b/src/srht_contrib/services/srht_client.py
@@ -0,0 +1,77 @@
+from __future__ import annotations
+
+import logging
+from typing import Any
+
+import httpx
+
+
+logger = logging.getLogger(__name__)
+
+
+class SourceHutClientError(RuntimeError):
+ """Raised when a SourceHut GraphQL request fails."""
+
+
+class SourceHutGraphQLClient:
+ def __init__(
+ self,
+ endpoint: str,
+ token: str,
+ *,
+ timeout: float = 15.0,
+ max_retries: int = 2,
+ transport: httpx.BaseTransport | None = None,
+ ) -> None:
+ self.endpoint = endpoint
+ self.timeout = timeout
+ self.max_retries = max_retries
+ headers = {
+ "Authorization": f"Bearer {token}",
+ "Content-Type": "application/json",
+ }
+ self._client = httpx.Client(headers=headers, timeout=timeout, transport=transport)
+
+ def execute(self, query: str, variables: dict[str, Any] | None = None) -> dict[str, Any]:
+ payload = {"query": query, "variables": variables or {}}
+ attempts = self.max_retries + 1
+
+ for attempt in range(1, attempts + 1):
+ try:
+ response = self._client.post(self.endpoint, json=payload)
+ response.raise_for_status()
+ body = response.json()
+ except httpx.HTTPStatusError as exc:
+ response_text = exc.response.text[:500]
+ logger.warning(
+ "SourceHut HTTP failure from %s on attempt %s/%s: %s %s",
+ self.endpoint,
+ attempt,
+ attempts,
+ exc.response.status_code,
+ response_text,
+ )
+ if exc.response.status_code >= 500 and attempt < attempts:
+ continue
+ raise SourceHutClientError(
+ f"HTTP error from SourceHut: {exc.response.status_code} {response_text}".strip()
+ ) from exc
+ except httpx.HTTPError as exc:
+ logger.warning(
+ "SourceHut network failure from %s on attempt %s/%s",
+ self.endpoint,
+ attempt,
+ attempts,
+ )
+ if attempt < attempts:
+ continue
+ raise SourceHutClientError("Network error while contacting SourceHut") from exc
+
+ if "errors" in body:
+ raise SourceHutClientError(f"GraphQL errors returned by SourceHut: {body['errors']}")
+ return body.get("data", {})
+
+ raise SourceHutClientError("SourceHut request exhausted retries")
+
+ def close(self) -> None:
+ self._client.close()
diff --git a/src/srht_contrib/services/todo.py b/src/srht_contrib/services/todo.py
new file mode 100644
index 0000000..2d25acf
--- /dev/null
+++ b/src/srht_contrib/services/todo.py
@@ -0,0 +1,559 @@
+from __future__ import annotations
+
+from dataclasses import dataclass
+from datetime import UTC, datetime, timedelta
+import logging
+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.utils.dates import ensure_utc, parse_datetime
+
+
+logger = logging.getLogger(__name__)
+
+
+TODO_ACTIVITY_QUERY = """
+query TodoActivity($cursor: Cursor) {
+ me {
+ canonicalName
+ }
+ events(cursor: $cursor) {
+ results {
+ id
+ created
+ ticket {
+ id
+ ref
+ status
+ resolution
+ tracker {
+ name
+ }
+ }
+ changes {
+ __typename
+ eventType
+ ticket {
+ id
+ }
+ ... on Created {
+ author {
+ canonicalName
+ }
+ }
+ ... on Comment {
+ author {
+ canonicalName
+ }
+ }
+ ... on StatusChange {
+ editor {
+ canonicalName
+ }
+ oldStatus
+ newStatus
+ oldResolution
+ newResolution
+ }
+ }
+ }
+ cursor
+ }
+}
+""".strip()
+
+TODO_TRACKERS_QUERY = """
+query TodoTrackers($cursor: Cursor) {
+ me {
+ canonicalName
+ trackers(cursor: $cursor) {
+ results {
+ id
+ rid
+ name
+ }
+ cursor
+ }
+ }
+}
+""".strip()
+
+TODO_TRACKER_TICKETS_QUERY = """
+query TodoTrackerTickets($trackerRid: ID!, $cursor: Cursor) {
+ tracker(rid: $trackerRid) {
+ id
+ name
+ tickets(cursor: $cursor) {
+ results {
+ id
+ ref
+ created
+ updated
+ status
+ resolution
+ submitter {
+ canonicalName
+ }
+ }
+ cursor
+ }
+ }
+}
+""".strip()
+
+TODO_TICKET_EVENTS_QUERY = """
+query TodoTicketEvents($trackerRid: ID!, $ticketId: Int!, $cursor: Cursor) {
+ tracker(rid: $trackerRid) {
+ id
+ name
+ ticket(id: $ticketId) {
+ id
+ ref
+ status
+ resolution
+ events(cursor: $cursor) {
+ results {
+ id
+ created
+ changes {
+ __typename
+ eventType
+ ticket {
+ id
+ }
+ ... on Created {
+ author {
+ canonicalName
+ }
+ }
+ ... on Comment {
+ author {
+ canonicalName
+ }
+ }
+ ... on StatusChange {
+ editor {
+ canonicalName
+ }
+ oldStatus
+ newStatus
+ oldResolution
+ newResolution
+ }
+ }
+ }
+ cursor
+ }
+ }
+ }
+}
+""".strip()
+
+
+TICKET_CLOSED_STATUSES = {"RESOLVED"}
+TICKET_CLOSED_RESOLUTIONS = {
+ "CLOSED",
+ "FIXED",
+ "IMPLEMENTED",
+ "WONT_FIX",
+ "BY_DESIGN",
+ "INVALID",
+ "DUPLICATE",
+ "NOT_OUR_BUG",
+}
+
+
+class TodoSchemaError(RuntimeError):
+ """Raised when SourceHut returns an unexpected todo event shape."""
+
+
+def _safe_nested_name(entity: dict[str, Any] | None) -> str | None:
+ if not entity:
+ return None
+ return entity.get("canonicalName") or entity.get("name")
+
+
+def _repo_name_from_event(event: dict[str, Any]) -> str | None:
+ ticket = event.get("ticket") or {}
+ tracker = ticket.get("tracker") or {}
+ return tracker.get("name")
+
+
+def _resource_id_from_event(event: dict[str, Any]) -> str:
+ ticket = event.get("ticket") or {}
+ return str(ticket.get("ref") or ticket.get("id") or event["id"])
+
+
+def _change_ticket_id(change: dict[str, Any], event: dict[str, Any]) -> str:
+ ticket = change.get("ticket") or event.get("ticket") or {}
+ return str(ticket.get("id") or event["id"])
+
+
+def _normalize_event_change(
+ *,
+ settings: Settings,
+ actor: str,
+ event: dict[str, Any],
+ change: dict[str, Any],
+ occurred_at: datetime,
+) -> NormalizedEvent | None:
+ change_type = change.get("__typename")
+ event_id = str(event["id"])
+ resource_id = _resource_id_from_event(event)
+ repo_name = _repo_name_from_event(event)
+ ticket_id = _change_ticket_id(change, event)
+
+ if change_type == "Created" and _safe_nested_name(change.get("author")) == actor:
+ return NormalizedEvent(
+ service="todo",
+ event_type="ticket_created",
+ actor=actor,
+ repo_name=repo_name,
+ resource_id=resource_id,
+ external_uid=f"todo:event:{event_id}:created:{ticket_id}",
+ occurred_at=occurred_at,
+ weight=settings.event_weights["ticket_created"],
+ raw_payload_json={"event": event, "change": change},
+ )
+
+ if change_type == "Comment" and _safe_nested_name(change.get("author")) == actor:
+ return NormalizedEvent(
+ service="todo",
+ event_type="ticket_comment",
+ actor=actor,
+ repo_name=repo_name,
+ resource_id=resource_id,
+ external_uid=f"todo:event:{event_id}:comment:{ticket_id}",
+ occurred_at=occurred_at,
+ weight=settings.event_weights["ticket_comment"],
+ raw_payload_json={"event": event, "change": change},
+ )
+
+ if change_type == "StatusChange" and _safe_nested_name(change.get("editor")) == actor:
+ new_status = change.get("newStatus")
+ new_resolution = change.get("newResolution")
+ if new_status in TICKET_CLOSED_STATUSES or new_resolution in TICKET_CLOSED_RESOLUTIONS:
+ return NormalizedEvent(
+ service="todo",
+ event_type="ticket_closed",
+ actor=actor,
+ repo_name=repo_name,
+ resource_id=resource_id,
+ external_uid=f"todo:event:{event_id}:closed:{ticket_id}",
+ occurred_at=occurred_at,
+ weight=settings.event_weights["ticket_closed"],
+ raw_payload_json={"event": event, "change": change},
+ )
+
+ return None
+
+
+def _extract_event_cursor_page(data: dict[str, Any]) -> tuple[str | None, list[dict[str, Any]], str]:
+ me = data.get("me") or {}
+ canonical_actor = me.get("canonicalName")
+ if not canonical_actor:
+ raise TodoSchemaError("todo.sr.ht response did not include me.canonicalName")
+
+ events = data.get("events") or {}
+ results = events.get("results") or []
+ if not isinstance(results, list):
+ raise TodoSchemaError("todo.sr.ht response did not include events.results")
+
+ return events.get("cursor"), [event for event in results if isinstance(event, dict)], canonical_actor
+
+
+@dataclass(slots=True)
+class TodoPollResult:
+ events: list[NormalizedEvent]
+ cursor: str
+
+
+class TodoIngestionService:
+ """Fetches todo.sr.ht activity from the authenticated event feed and normalizes it."""
+
+ service_name = "todo"
+
+ def __init__(self, client: SourceHutGraphQLClient, settings: Settings) -> None:
+ self.client = client
+ self.settings = settings
+
+ def fetch_recent_events(self, actor: str, since: datetime | None = None) -> TodoPollResult:
+ since_dt = ensure_utc(since or (datetime.now(tz=UTC) - timedelta(days=30)))
+ cursor_time = datetime.now(tz=UTC).isoformat()
+ feed_result = self._fetch_from_activity_feed(actor=actor, since=since_dt)
+ if feed_result.events:
+ return TodoPollResult(events=feed_result.events, cursor=cursor_time)
+
+ logger.info(
+ "todo activity feed returned no normalized events for actor=%s; falling back to tracker crawl",
+ actor,
+ )
+ tracker_events = self._fetch_from_trackers(actor=actor, since=since_dt)
+ return TodoPollResult(events=tracker_events, cursor=cursor_time)
+
+ def _fetch_from_activity_feed(self, actor: str, since: datetime) -> TodoPollResult:
+ events: list[NormalizedEvent] = []
+ cursor: str | None = None
+ effective_actor = actor
+
+ for _ in range(10):
+ data = self.client.execute(TODO_ACTIVITY_QUERY, {"cursor": cursor})
+ cursor, page_events, canonical_actor = _extract_event_cursor_page(data)
+ effective_actor = actor or canonical_actor
+ logger.info(
+ "todo page fetched for actor=%s canonical_actor=%s events=%s next_cursor=%s since=%s",
+ actor,
+ canonical_actor,
+ len(page_events),
+ bool(cursor),
+ since.isoformat(),
+ )
+
+ stop_paging = False
+ for event in page_events:
+ occurred_at = parse_datetime(event["created"])
+ event_id = str(event.get("id"))
+ resource_id = _resource_id_from_event(event)
+ repo_name = _repo_name_from_event(event)
+ change_list = event.get("changes") or []
+ logger.info(
+ "todo event id=%s resource=%s repo=%s occurred_at=%s changes=%s",
+ event_id,
+ resource_id,
+ repo_name,
+ occurred_at.isoformat(),
+ len(change_list) if isinstance(change_list, list) else "unknown",
+ )
+ if occurred_at < since:
+ stop_paging = True
+ logger.info(
+ "todo event id=%s skipped because occurred_at=%s is before since=%s",
+ event_id,
+ occurred_at.isoformat(),
+ since.isoformat(),
+ )
+ continue
+
+ for change in change_list:
+ if not isinstance(change, dict):
+ logger.info("todo event id=%s skipped non-dict change payload", event_id)
+ continue
+ change_type = change.get("__typename")
+ author = _safe_nested_name(change.get("author"))
+ editor = _safe_nested_name(change.get("editor"))
+ logger.info(
+ "todo change event_id=%s type=%s eventType=%s author=%s editor=%s newStatus=%s newResolution=%s",
+ event_id,
+ change_type,
+ change.get("eventType"),
+ author,
+ editor,
+ change.get("newStatus"),
+ change.get("newResolution"),
+ )
+ normalized = _normalize_event_change(
+ settings=self.settings,
+ actor=effective_actor,
+ event=event,
+ change=change,
+ occurred_at=occurred_at,
+ )
+ if normalized is not None:
+ logger.info(
+ "todo change accepted event_id=%s normalized_type=%s external_uid=%s",
+ event_id,
+ normalized.event_type,
+ normalized.external_uid,
+ )
+ events.append(normalized)
+ else:
+ logger.info(
+ "todo change skipped event_id=%s for actor=%s",
+ event_id,
+ effective_actor,
+ )
+
+ if stop_paging or not cursor:
+ break
+
+ logger.info("todo poll complete for actor=%s normalized_events=%s", effective_actor, len(events))
+ return TodoPollResult(events=events, cursor=datetime.now(tz=UTC).isoformat())
+
+ def _fetch_from_trackers(self, actor: str, since: datetime) -> list[NormalizedEvent]:
+ data = self.client.execute(TODO_TRACKERS_QUERY, {"cursor": None})
+ me = data.get("me") or {}
+ canonical_actor = me.get("canonicalName") or actor
+ trackers = ((me.get("trackers") or {}).get("results")) or []
+ logger.info("todo tracker crawl actor=%s trackers=%s", canonical_actor, len(trackers))
+
+ events: list[NormalizedEvent] = []
+ for tracker in trackers:
+ if not isinstance(tracker, dict):
+ continue
+ tracker_id = tracker.get("id")
+ tracker_rid = tracker.get("rid")
+ tracker_name = tracker.get("name")
+ tracker_events = self._fetch_tracker_tickets(
+ actor=canonical_actor,
+ tracker_id=str(tracker_id),
+ tracker_rid=str(tracker_rid),
+ tracker_name=tracker_name,
+ since=since,
+ )
+ events.extend(tracker_events)
+
+ logger.info("todo tracker crawl complete actor=%s normalized_events=%s", canonical_actor, len(events))
+ return events
+
+ def _fetch_tracker_tickets(
+ self,
+ *,
+ actor: str,
+ tracker_id: str,
+ tracker_rid: str,
+ tracker_name: str | None,
+ since: datetime,
+ ) -> list[NormalizedEvent]:
+ events: list[NormalizedEvent] = []
+ cursor: str | None = None
+
+ for _ in range(50):
+ data = self.client.execute(
+ TODO_TRACKER_TICKETS_QUERY,
+ {"trackerRid": tracker_rid, "cursor": cursor},
+ )
+ tracker = data.get("tracker") or {}
+ tickets_page = tracker.get("tickets") or {}
+ tickets = tickets_page.get("results") or []
+ cursor = tickets_page.get("cursor")
+ logger.info(
+ "todo tracker=%s ticket page count=%s next_cursor=%s",
+ tracker_name or tracker_id,
+ len(tickets),
+ bool(cursor),
+ )
+
+ stop_paging = False
+ for ticket in tickets:
+ if not isinstance(ticket, dict):
+ continue
+ updated_at = parse_datetime(ticket["updated"])
+ if updated_at < since:
+ stop_paging = True
+ logger.info(
+ "todo ticket ref=%s skipped because updated_at=%s is before since=%s",
+ ticket.get("ref"),
+ updated_at.isoformat(),
+ since.isoformat(),
+ )
+ continue
+ events.extend(
+ self._fetch_ticket_events(
+ actor=actor,
+ tracker_id=tracker_id,
+ tracker_rid=tracker_rid,
+ tracker_name=tracker_name,
+ ticket=ticket,
+ since=since,
+ )
+ )
+
+ if stop_paging or not cursor:
+ break
+
+ return events
+
+ def _fetch_ticket_events(
+ self,
+ *,
+ actor: str,
+ tracker_id: str,
+ tracker_rid: str,
+ tracker_name: str | None,
+ ticket: dict[str, Any],
+ since: datetime,
+ ) -> list[NormalizedEvent]:
+ events: list[NormalizedEvent] = []
+ cursor: str | None = None
+ ticket_id = int(ticket["id"])
+ ticket_ref = str(ticket.get("ref") or ticket_id)
+
+ for _ in range(50):
+ data = self.client.execute(
+ TODO_TICKET_EVENTS_QUERY,
+ {"trackerRid": tracker_rid, "ticketId": ticket_id, "cursor": 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 []
+ cursor = event_page.get("cursor")
+ logger.info(
+ "todo ticket events ref=%s tracker=%s count=%s next_cursor=%s",
+ ticket_ref,
+ tracker_name or tracker_id,
+ len(page_events),
+ bool(cursor),
+ )
+
+ stop_paging = False
+ for event in page_events:
+ if not isinstance(event, dict):
+ continue
+ event["ticket"] = {
+ "id": ticket_payload.get("id", ticket.get("id")),
+ "ref": ticket_payload.get("ref", ticket_ref),
+ "status": ticket_payload.get("status", ticket.get("status")),
+ "resolution": ticket_payload.get("resolution", ticket.get("resolution")),
+ "tracker": {"name": tracker_name},
+ }
+ occurred_at = parse_datetime(event["created"])
+ event_id = str(event.get("id"))
+ change_list = event.get("changes") or []
+ logger.info(
+ "todo ticket event ref=%s id=%s occurred_at=%s changes=%s",
+ ticket_ref,
+ event_id,
+ occurred_at.isoformat(),
+ len(change_list) if isinstance(change_list, list) else "unknown",
+ )
+ if occurred_at < since:
+ stop_paging = True
+ continue
+
+ for change in change_list:
+ if not isinstance(change, dict):
+ continue
+ logger.info(
+ "todo ticket change ref=%s event_id=%s type=%s eventType=%s author=%s editor=%s newStatus=%s newResolution=%s",
+ ticket_ref,
+ event_id,
+ change.get("__typename"),
+ change.get("eventType"),
+ _safe_nested_name(change.get("author")),
+ _safe_nested_name(change.get("editor")),
+ change.get("newStatus"),
+ change.get("newResolution"),
+ )
+ normalized = _normalize_event_change(
+ settings=self.settings,
+ actor=actor,
+ event=event,
+ change=change,
+ occurred_at=occurred_at,
+ )
+ if normalized is not None:
+ logger.info(
+ "todo ticket change accepted ref=%s event_id=%s normalized_type=%s",
+ ticket_ref,
+ event_id,
+ normalized.event_type,
+ )
+ events.append(normalized)
+
+ if stop_paging or not cursor:
+ break
+
+ return events