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