krz/hutch-stats

Server-side utility for calculating contributions for sourcehut users.

clone: git clone https://gitbay.org/krz/hutch-stats.git

main: src/srht_contrib/services/todo.py · raw

  1from __future__ import annotations
  2
  3import copy
  4from dataclasses import dataclass
  5from datetime import UTC, datetime, timedelta
  6import logging
  7from typing import Any
  8
  9from srht_contrib.config import Settings
 10from srht_contrib.schemas import NormalizedEvent
 11from srht_contrib.services.srht_client import SourceHutGraphQLClient
 12from srht_contrib.services.types import BackfillBatchResult
 13from srht_contrib.utils.dates import ensure_utc, parse_datetime
 14
 15
 16logger = logging.getLogger(__name__)
 17
 18
 19TODO_ACTIVITY_QUERY = """
 20query TodoActivity($cursor: Cursor) {
 21  me {
 22    canonicalName
 23  }
 24  events(cursor: $cursor) {
 25    results {
 26      id
 27      created
 28      ticket {
 29        id
 30        ref
 31        status
 32        resolution
 33        tracker {
 34          name
 35        }
 36      }
 37      changes {
 38        __typename
 39        eventType
 40        ticket {
 41          id
 42        }
 43        ... on Created {
 44          author {
 45            canonicalName
 46          }
 47        }
 48        ... on Comment {
 49          author {
 50            canonicalName
 51          }
 52        }
 53        ... on StatusChange {
 54          editor {
 55            canonicalName
 56          }
 57          oldStatus
 58          newStatus
 59          oldResolution
 60          newResolution
 61        }
 62      }
 63    }
 64    cursor
 65  }
 66}
 67""".strip()
 68
 69TODO_TRACKERS_QUERY = """
 70query TodoTrackers($cursor: Cursor) {
 71  me {
 72    canonicalName
 73    trackers(cursor: $cursor) {
 74      results {
 75        id
 76        rid
 77        name
 78      }
 79      cursor
 80    }
 81  }
 82}
 83""".strip()
 84
 85TODO_TRACKER_TICKETS_QUERY = """
 86query TodoTrackerTickets($trackerRid: ID!, $cursor: Cursor) {
 87  tracker(rid: $trackerRid) {
 88    id
 89    name
 90    tickets(cursor: $cursor) {
 91      results {
 92        id
 93        ref
 94        created
 95        updated
 96        status
 97        resolution
 98        submitter {
 99          canonicalName
100        }
101      }
102      cursor
103    }
104  }
105}
106""".strip()
107
108TODO_TICKET_EVENTS_QUERY = """
109query TodoTicketEvents($trackerRid: ID!, $ticketId: Int!, $cursor: Cursor) {
110  tracker(rid: $trackerRid) {
111    id
112    name
113    ticket(id: $ticketId) {
114      id
115      ref
116      status
117      resolution
118      events(cursor: $cursor) {
119        results {
120          id
121          created
122          changes {
123            __typename
124            eventType
125            ticket {
126              id
127            }
128            ... on Created {
129              author {
130                canonicalName
131              }
132            }
133            ... on Comment {
134              author {
135                canonicalName
136              }
137            }
138            ... on StatusChange {
139              editor {
140                canonicalName
141              }
142              oldStatus
143              newStatus
144              oldResolution
145              newResolution
146            }
147          }
148        }
149        cursor
150      }
151    }
152  }
153}
154""".strip()
155
156
157TICKET_CLOSED_STATUSES = {"RESOLVED"}
158TICKET_CLOSED_RESOLUTIONS = {
159    "CLOSED",
160    "FIXED",
161    "IMPLEMENTED",
162    "WONT_FIX",
163    "BY_DESIGN",
164    "INVALID",
165    "DUPLICATE",
166    "NOT_OUR_BUG",
167}
168
169
170class TodoSchemaError(RuntimeError):
171    """Raised when SourceHut returns an unexpected todo event shape."""
172
173
174def _safe_nested_name(entity: dict[str, Any] | None) -> str | None:
175    if not entity:
176        return None
177    return entity.get("canonicalName") or entity.get("name")
178
179
180def _repo_name_from_event(event: dict[str, Any]) -> str | None:
181    ticket = event.get("ticket") or {}
182    tracker = ticket.get("tracker") or {}
183    return tracker.get("name")
184
185
186def _resource_id_from_event(event: dict[str, Any]) -> str:
187    ticket = event.get("ticket") or {}
188    return str(ticket.get("ref") or ticket.get("id") or event["id"])
189
190
191def _change_ticket_id(change: dict[str, Any], event: dict[str, Any]) -> str:
192    ticket = change.get("ticket") or event.get("ticket") or {}
193    return str(ticket.get("id") or event["id"])
194
195
196def _normalize_event_change(
197    *,
198    settings: Settings,
199    actor: str,
200    event: dict[str, Any],
201    change: dict[str, Any],
202    occurred_at: datetime,
203) -> NormalizedEvent | None:
204    change_type = change.get("__typename")
205    event_id = str(event["id"])
206    resource_id = _resource_id_from_event(event)
207    repo_name = _repo_name_from_event(event)
208    ticket_id = _change_ticket_id(change, event)
209
210    if change_type == "Created" and _safe_nested_name(change.get("author")) == actor:
211        return NormalizedEvent(
212            service="todo",
213            event_type="ticket_created",
214            actor=actor,
215            repo_name=repo_name,
216            resource_id=resource_id,
217            external_uid=f"todo:event:{event_id}:created:{ticket_id}",
218            occurred_at=occurred_at,
219            weight=settings.event_weights["ticket_created"],
220            raw_payload_json={"event": event, "change": change},
221        )
222
223    if change_type == "Comment" and _safe_nested_name(change.get("author")) == actor:
224        return NormalizedEvent(
225            service="todo",
226            event_type="ticket_comment",
227            actor=actor,
228            repo_name=repo_name,
229            resource_id=resource_id,
230            external_uid=f"todo:event:{event_id}:comment:{ticket_id}",
231            occurred_at=occurred_at,
232            weight=settings.event_weights["ticket_comment"],
233            raw_payload_json={"event": event, "change": change},
234        )
235
236    if change_type == "StatusChange" and _safe_nested_name(change.get("editor")) == actor:
237        new_status = change.get("newStatus")
238        new_resolution = change.get("newResolution")
239        if new_status in TICKET_CLOSED_STATUSES or new_resolution in TICKET_CLOSED_RESOLUTIONS:
240            return NormalizedEvent(
241                service="todo",
242                event_type="ticket_closed",
243                actor=actor,
244                repo_name=repo_name,
245                resource_id=resource_id,
246                external_uid=f"todo:event:{event_id}:closed:{ticket_id}",
247                occurred_at=occurred_at,
248                weight=settings.event_weights["ticket_closed"],
249                raw_payload_json={"event": event, "change": change},
250            )
251
252    return None
253
254
255def _extract_event_cursor_page(data: dict[str, Any]) -> tuple[str | None, list[dict[str, Any]], str]:
256    me = data.get("me") or {}
257    canonical_actor = me.get("canonicalName")
258    if not canonical_actor:
259        raise TodoSchemaError("todo.sr.ht response did not include me.canonicalName")
260
261    events = data.get("events") or {}
262    results = events.get("results") or []
263    if not isinstance(results, list):
264        raise TodoSchemaError("todo.sr.ht response did not include events.results")
265
266    return events.get("cursor"), [event for event in results if isinstance(event, dict)], canonical_actor
267
268
269@dataclass(slots=True)
270class TodoPollResult:
271    events: list[NormalizedEvent]
272    cursor: str
273
274
275class TodoIngestionService:
276    """Fetches todo.sr.ht activity from the authenticated event feed and normalizes it."""
277
278    service_name = "todo"
279
280    def __init__(self, client: SourceHutGraphQLClient, settings: Settings) -> None:
281        self.client = client
282        self.settings = settings
283
284    def fetch_recent_events(self, actor: str, since: datetime | None = None) -> TodoPollResult:
285        since_dt = ensure_utc(since or (datetime.now(tz=UTC) - timedelta(days=30)))
286        cursor_time = datetime.now(tz=UTC).isoformat()
287        feed_result = self._fetch_from_activity_feed(actor=actor, since=since_dt)
288        if feed_result.events:
289            return TodoPollResult(events=feed_result.events, cursor=cursor_time)
290
291        logger.info(
292            "todo activity feed returned no normalized events for actor=%s; falling back to tracker crawl",
293            actor,
294        )
295        tracker_events = self._fetch_from_trackers(actor=actor, since=since_dt)
296        return TodoPollResult(events=tracker_events, cursor=cursor_time)
297
298    def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
299        return self._fetch_backfill_batch(actor, cursor_state, since=None)
300
301    def fetch_recent_backfill_batch(
302        self,
303        actor: str,
304        cursor_state: dict | None = None,
305        *,
306        since: datetime,
307    ) -> BackfillBatchResult:
308        return self._fetch_backfill_batch(actor, cursor_state, since=since)
309
310    def _fetch_backfill_batch(
311        self,
312        actor: str,
313        cursor_state: dict | None,
314        *,
315        since: datetime | None,
316    ) -> BackfillBatchResult:
317        state = {
318            "trackers_cursor": None,
319            "tracker_queue": [],
320            "current_tracker": None,
321            "current_ticket": None,
322            "trackers_loaded": False,
323        }
324        if cursor_state:
325            state.update(copy.deepcopy(cursor_state))
326
327        if not state["trackers_loaded"]:
328            data = self.client.execute(TODO_TRACKERS_QUERY, {"cursor": state["trackers_cursor"]})
329            me = data.get("me") or {}
330            trackers_page = me.get("trackers") or {}
331            results = trackers_page.get("results") or []
332            state["trackers_cursor"] = trackers_page.get("cursor")
333            for tracker in results:
334                if not isinstance(tracker, dict):
335                    continue
336                state["tracker_queue"].append(
337                    {
338                        "id": str(tracker.get("id")),
339                        "rid": str(tracker.get("rid")),
340                        "name": tracker.get("name"),
341                        "tickets_cursor": None,
342                        "pending_tickets": [],
343                    }
344                )
345            state["trackers_loaded"] = not bool(state["trackers_cursor"])
346            logger.info(
347                "todo backfill tracker discovery actor=%s trackers=%s next_cursor=%s queue=%s",
348                actor,
349                len(results),
350                bool(state["trackers_cursor"]),
351                len(state["tracker_queue"]),
352            )
353            complete = state["trackers_loaded"] and not state["tracker_queue"]
354            return BackfillBatchResult(events=[], cursor_state=None if complete else state, complete=complete)
355
356        if state["current_tracker"] is None:
357            if not state["tracker_queue"]:
358                return BackfillBatchResult(events=[], cursor_state=None, complete=True)
359            state["current_tracker"] = state["tracker_queue"].pop(0)
360
361        current_tracker = state["current_tracker"]
362        if state["current_ticket"] is None and not current_tracker["pending_tickets"]:
363            data = self.client.execute(
364                TODO_TRACKER_TICKETS_QUERY,
365                {"trackerRid": current_tracker["rid"], "cursor": current_tracker["tickets_cursor"]},
366            )
367            tracker = data.get("tracker") or {}
368            tickets_page = tracker.get("tickets") or {}
369            tickets = tickets_page.get("results") or []
370            next_tickets_cursor = tickets_page.get("cursor")
371            stop_tracker_paging = False
372            for ticket in tickets:
373                if not isinstance(ticket, dict):
374                    continue
375                if since is not None:
376                    updated = parse_datetime(ticket["updated"])
377                    if updated < since:
378                        stop_tracker_paging = True
379                        continue
380                current_tracker["pending_tickets"].append(
381                    {"id": int(ticket["id"]), "ref": str(ticket.get("ref") or ticket["id"])}
382                )
383            current_tracker["tickets_cursor"] = None if stop_tracker_paging else next_tickets_cursor
384            logger.info(
385                "todo backfill tracker=%s ticket page count=%s next_cursor=%s pending_tickets=%s",
386                current_tracker.get("name") or current_tracker.get("id"),
387                len(tickets),
388                bool(current_tracker["tickets_cursor"]),
389                len(current_tracker["pending_tickets"]),
390            )
391            if not current_tracker["pending_tickets"] and not current_tracker["tickets_cursor"]:
392                state["current_tracker"] = None
393            return BackfillBatchResult(events=[], cursor_state=state, complete=False)
394
395        if state["current_ticket"] is None:
396            if current_tracker["pending_tickets"]:
397                next_ticket = current_tracker["pending_tickets"].pop(0)
398                state["current_ticket"] = {**next_ticket, "cursor": None}
399            else:
400                state["current_tracker"] = None
401                return BackfillBatchResult(events=[], cursor_state=state, complete=False)
402
403        current_ticket = state["current_ticket"]
404        data = self.client.execute(
405            TODO_TICKET_EVENTS_QUERY,
406            {"trackerRid": current_tracker["rid"], "ticketId": current_ticket["id"], "cursor": current_ticket["cursor"]},
407        )
408        tracker = data.get("tracker") or {}
409        ticket_payload = (tracker.get("ticket") or {}) if isinstance(tracker, dict) else {}
410        event_page = ticket_payload.get("events") or {}
411        page_events = event_page.get("results") or []
412        next_cursor = event_page.get("cursor")
413        logger.info(
414            "todo backfill ticket events ref=%s tracker=%s count=%s next_cursor=%s",
415            current_ticket["ref"],
416            current_tracker.get("name") or current_tracker.get("id"),
417            len(page_events),
418            bool(next_cursor),
419        )
420
421        events: list[NormalizedEvent] = []
422        stop_ticket_paging = False
423        for event in page_events:
424            if not isinstance(event, dict):
425                continue
426            event["ticket"] = {
427                "id": ticket_payload.get("id", current_ticket["id"]),
428                "ref": ticket_payload.get("ref", current_ticket["ref"]),
429                "status": ticket_payload.get("status"),
430                "resolution": ticket_payload.get("resolution"),
431                "tracker": {"name": current_tracker.get("name")},
432            }
433            occurred_at = parse_datetime(event["created"])
434            if since is not None and occurred_at < since:
435                stop_ticket_paging = True
436                continue
437            for change in event.get("changes") or []:
438                if not isinstance(change, dict):
439                    continue
440                normalized = _normalize_event_change(
441                    settings=self.settings,
442                    actor=actor,
443                    event=event,
444                    change=change,
445                    occurred_at=occurred_at,
446                )
447                if normalized is not None:
448                    events.append(normalized)
449
450        if next_cursor and not stop_ticket_paging:
451            state["current_ticket"]["cursor"] = next_cursor
452        else:
453            state["current_ticket"] = None
454            if not current_tracker["pending_tickets"] and not current_tracker["tickets_cursor"]:
455                state["current_tracker"] = None
456
457        complete = (
458            state["trackers_loaded"]
459            and not state["tracker_queue"]
460            and state["current_tracker"] is None
461            and state["current_ticket"] is None
462        )
463        return BackfillBatchResult(events=events, cursor_state=None if complete else state, complete=complete)
464
465    def _fetch_from_activity_feed(self, actor: str, since: datetime) -> TodoPollResult:
466        events: list[NormalizedEvent] = []
467        cursor: str | None = None
468        effective_actor = actor
469
470        for _ in range(10):
471            data = self.client.execute(TODO_ACTIVITY_QUERY, {"cursor": cursor})
472            cursor, page_events, canonical_actor = _extract_event_cursor_page(data)
473            effective_actor = actor or canonical_actor
474            logger.info(
475                "todo page fetched for actor=%s canonical_actor=%s events=%s next_cursor=%s since=%s",
476                actor,
477                canonical_actor,
478                len(page_events),
479                bool(cursor),
480                since.isoformat(),
481            )
482
483            stop_paging = False
484            for event in page_events:
485                occurred_at = parse_datetime(event["created"])
486                event_id = str(event.get("id"))
487                resource_id = _resource_id_from_event(event)
488                repo_name = _repo_name_from_event(event)
489                change_list = event.get("changes") or []
490                logger.info(
491                    "todo event id=%s resource=%s repo=%s occurred_at=%s changes=%s",
492                    event_id,
493                    resource_id,
494                    repo_name,
495                    occurred_at.isoformat(),
496                    len(change_list) if isinstance(change_list, list) else "unknown",
497                )
498                if occurred_at < since:
499                    stop_paging = True
500                    logger.info(
501                        "todo event id=%s skipped because occurred_at=%s is before since=%s",
502                        event_id,
503                        occurred_at.isoformat(),
504                        since.isoformat(),
505                    )
506                    continue
507
508                for change in change_list:
509                    if not isinstance(change, dict):
510                        logger.info("todo event id=%s skipped non-dict change payload", event_id)
511                        continue
512                    change_type = change.get("__typename")
513                    author = _safe_nested_name(change.get("author"))
514                    editor = _safe_nested_name(change.get("editor"))
515                    logger.info(
516                        "todo change event_id=%s type=%s eventType=%s author=%s editor=%s newStatus=%s newResolution=%s",
517                        event_id,
518                        change_type,
519                        change.get("eventType"),
520                        author,
521                        editor,
522                        change.get("newStatus"),
523                        change.get("newResolution"),
524                    )
525                    normalized = _normalize_event_change(
526                        settings=self.settings,
527                        actor=effective_actor,
528                        event=event,
529                        change=change,
530                        occurred_at=occurred_at,
531                    )
532                    if normalized is not None:
533                        logger.info(
534                            "todo change accepted event_id=%s normalized_type=%s external_uid=%s",
535                            event_id,
536                            normalized.event_type,
537                            normalized.external_uid,
538                        )
539                        events.append(normalized)
540                    else:
541                        logger.info(
542                            "todo change skipped event_id=%s for actor=%s",
543                            event_id,
544                            effective_actor,
545                        )
546
547            if stop_paging or not cursor:
548                break
549
550        logger.info("todo poll complete for actor=%s normalized_events=%s", effective_actor, len(events))
551        return TodoPollResult(events=events, cursor=datetime.now(tz=UTC).isoformat())
552
553    def _fetch_from_trackers(self, actor: str, since: datetime) -> list[NormalizedEvent]:
554        data = self.client.execute(TODO_TRACKERS_QUERY, {"cursor": None})
555        me = data.get("me") or {}
556        canonical_actor = me.get("canonicalName") or actor
557        trackers = ((me.get("trackers") or {}).get("results")) or []
558        logger.info("todo tracker crawl actor=%s trackers=%s", canonical_actor, len(trackers))
559
560        events: list[NormalizedEvent] = []
561        for tracker in trackers:
562            if not isinstance(tracker, dict):
563                continue
564            tracker_id = tracker.get("id")
565            tracker_rid = tracker.get("rid")
566            tracker_name = tracker.get("name")
567            tracker_events = self._fetch_tracker_tickets(
568                actor=canonical_actor,
569                tracker_id=str(tracker_id),
570                tracker_rid=str(tracker_rid),
571                tracker_name=tracker_name,
572                since=since,
573            )
574            events.extend(tracker_events)
575
576        logger.info("todo tracker crawl complete actor=%s normalized_events=%s", canonical_actor, len(events))
577        return events
578
579    def _fetch_tracker_tickets(
580        self,
581        *,
582        actor: str,
583        tracker_id: str,
584        tracker_rid: str,
585        tracker_name: str | None,
586        since: datetime,
587    ) -> list[NormalizedEvent]:
588        events: list[NormalizedEvent] = []
589        cursor: str | None = None
590
591        for _ in range(50):
592            data = self.client.execute(
593                TODO_TRACKER_TICKETS_QUERY,
594                {"trackerRid": tracker_rid, "cursor": cursor},
595            )
596            tracker = data.get("tracker") or {}
597            tickets_page = tracker.get("tickets") or {}
598            tickets = tickets_page.get("results") or []
599            cursor = tickets_page.get("cursor")
600            logger.info(
601                "todo tracker=%s ticket page count=%s next_cursor=%s",
602                tracker_name or tracker_id,
603                len(tickets),
604                bool(cursor),
605            )
606
607            stop_paging = False
608            for ticket in tickets:
609                if not isinstance(ticket, dict):
610                    continue
611                updated_at = parse_datetime(ticket["updated"])
612                if updated_at < since:
613                    stop_paging = True
614                    logger.info(
615                        "todo ticket ref=%s skipped because updated_at=%s is before since=%s",
616                        ticket.get("ref"),
617                        updated_at.isoformat(),
618                        since.isoformat(),
619                    )
620                    continue
621                events.extend(
622                    self._fetch_ticket_events(
623                        actor=actor,
624                        tracker_id=tracker_id,
625                        tracker_rid=tracker_rid,
626                        tracker_name=tracker_name,
627                        ticket=ticket,
628                        since=since,
629                    )
630                )
631
632            if stop_paging or not cursor:
633                break
634
635        return events
636
637    def _fetch_ticket_events(
638        self,
639        *,
640        actor: str,
641        tracker_id: str,
642        tracker_rid: str,
643        tracker_name: str | None,
644        ticket: dict[str, Any],
645        since: datetime,
646    ) -> list[NormalizedEvent]:
647        events: list[NormalizedEvent] = []
648        cursor: str | None = None
649        ticket_id = int(ticket["id"])
650        ticket_ref = str(ticket.get("ref") or ticket_id)
651
652        for _ in range(50):
653            data = self.client.execute(
654                TODO_TICKET_EVENTS_QUERY,
655                {"trackerRid": tracker_rid, "ticketId": ticket_id, "cursor": cursor},
656            )
657            tracker = data.get("tracker") or {}
658            ticket_payload = (tracker.get("ticket") or {}) if isinstance(tracker, dict) else {}
659            event_page = ticket_payload.get("events") or {}
660            page_events = event_page.get("results") or []
661            cursor = event_page.get("cursor")
662            logger.info(
663                "todo ticket events ref=%s tracker=%s count=%s next_cursor=%s",
664                ticket_ref,
665                tracker_name or tracker_id,
666                len(page_events),
667                bool(cursor),
668            )
669
670            stop_paging = False
671            for event in page_events:
672                if not isinstance(event, dict):
673                    continue
674                event["ticket"] = {
675                    "id": ticket_payload.get("id", ticket.get("id")),
676                    "ref": ticket_payload.get("ref", ticket_ref),
677                    "status": ticket_payload.get("status", ticket.get("status")),
678                    "resolution": ticket_payload.get("resolution", ticket.get("resolution")),
679                    "tracker": {"name": tracker_name},
680                }
681                occurred_at = parse_datetime(event["created"])
682                event_id = str(event.get("id"))
683                change_list = event.get("changes") or []
684                logger.info(
685                    "todo ticket event ref=%s id=%s occurred_at=%s changes=%s",
686                    ticket_ref,
687                    event_id,
688                    occurred_at.isoformat(),
689                    len(change_list) if isinstance(change_list, list) else "unknown",
690                )
691                if occurred_at < since:
692                    stop_paging = True
693                    continue
694
695                for change in change_list:
696                    if not isinstance(change, dict):
697                        continue
698                    logger.info(
699                        "todo ticket change ref=%s event_id=%s type=%s eventType=%s author=%s editor=%s newStatus=%s newResolution=%s",
700                        ticket_ref,
701                        event_id,
702                        change.get("__typename"),
703                        change.get("eventType"),
704                        _safe_nested_name(change.get("author")),
705                        _safe_nested_name(change.get("editor")),
706                        change.get("newStatus"),
707                        change.get("newResolution"),
708                    )
709                    normalized = _normalize_event_change(
710                        settings=self.settings,
711                        actor=actor,
712                        event=event,
713                        change=change,
714                        occurred_at=occurred_at,
715                    )
716                    if normalized is not None:
717                        logger.info(
718                            "todo ticket change accepted ref=%s event_id=%s normalized_type=%s",
719                            ticket_ref,
720                            event_id,
721                            normalized.event_type,
722                        )
723                        events.append(normalized)
724
725            if stop_paging or not cursor:
726                break
727
728        return events