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/git.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 sqlalchemy import delete, select
10from sqlalchemy.orm import Session
11
12from srht_contrib.config import Settings
13from srht_contrib.models import DiscoveredRepository, TrackedRepository
14from srht_contrib.schemas import NormalizedEvent
15from srht_contrib.services.srht_client import SourceHutClientError, SourceHutGraphQLClient
16from srht_contrib.services.types import BackfillBatchResult
17from srht_contrib.utils.dates import ensure_utc, parse_datetime
18from srht_contrib.utils.identity import ActorIdentityResolver
19
20
21logger = logging.getLogger(__name__)
22
23
24REPOSITORY_BRANCHES_QUERY = """
25query RepositoryBranches($username: String!, $repoName: String!, $cursor: Cursor) {
26 user(username: $username) {
27 repository(name: $repoName) {
28 references(cursor: $cursor) {
29 results {
30 name
31 target
32 }
33 cursor
34 }
35 }
36 }
37}
38""".strip()
39
40
41REPOSITORY_LOG_QUERY = """
42query RepositoryLog($username: String!, $repoName: String!, $cursor: Cursor, $from: String) {
43 user(username: $username) {
44 repository(name: $repoName) {
45 name
46 owner {
47 canonicalName
48 }
49 log(cursor: $cursor, from: $from) {
50 results {
51 id
52 shortId
53 author {
54 name
55 email
56 time
57 }
58 committer {
59 name
60 email
61 time
62 }
63 message
64 }
65 cursor
66 }
67 }
68 }
69}
70""".strip()
71
72
73USER_REPOSITORIES_QUERY = """
74query UserRepositories($username: String!, $cursor: Cursor) {
75 user(username: $username) {
76 repositories(cursor: $cursor) {
77 results {
78 name
79 visibility
80 owner {
81 canonicalName
82 }
83 }
84 cursor
85 }
86 }
87}
88""".strip()
89
90
91@dataclass(slots=True)
92class GitPollResult:
93 events: list[NormalizedEvent]
94 cursor: str
95
96
97class GitIngestionService:
98 """Polls git.sr.ht repositories and normalizes commits for one actor."""
99
100 service_name = "git"
101
102 def __init__(self, client: SourceHutGraphQLClient, settings: Settings) -> None:
103 self.client = client
104 self.settings = settings
105 self.identity_resolver = ActorIdentityResolver(settings.actor_aliases_json)
106
107 def fetch_recent_events(
108 self,
109 actor: str,
110 since: datetime | None = None,
111 repositories: list[str] | None = None,
112 db: Session | None = None,
113 ) -> GitPollResult:
114 since_dt = ensure_utc(since or (datetime.now(tz=UTC) - timedelta(days=30)))
115 discovered_repositories = repositories or self._repositories_for_actor(actor, db)
116 if not discovered_repositories:
117 logger.info("git poll skipped for actor=%s because no repositories were discovered", actor)
118 return GitPollResult(events=[], cursor=datetime.now(tz=UTC).isoformat())
119
120 events: list[NormalizedEvent] = []
121 for repository in discovered_repositories:
122 owner, repo_name = self._split_repository(actor, repository)
123 try:
124 repo_events = self._fetch_repository_commits(actor=actor, owner=owner, repo_name=repo_name, since=since_dt)
125 except SourceHutClientError:
126 logger.warning("git poll skipped repository=%s/%s due to GraphQL error", owner, repo_name, exc_info=True)
127 continue
128 events.extend(repo_events)
129
130 logger.info("git poll complete for actor=%s normalized_events=%s", actor, len(events))
131 return GitPollResult(events=events, cursor=datetime.now(tz=UTC).isoformat())
132
133 def _repositories_for_actor(self, actor: str, db: Session | None = None) -> list[str]:
134 configured = {
135 self._canonical_repository_name(actor, repository)
136 for repository in self.settings.git_tracked_repositories
137 }
138 tracked = set()
139 discovered = set()
140 now = datetime.now(tz=UTC)
141
142 if db is not None:
143 tracked = set(
144 db.scalars(
145 select(TrackedRepository.repo_name)
146 .where(TrackedRepository.service == self.service_name)
147 .where(TrackedRepository.actor == actor)
148 ).all()
149 )
150 cutoff = now - timedelta(seconds=self.settings.git_repo_discovery_ttl_seconds)
151 discovered = set(
152 db.scalars(
153 select(DiscoveredRepository.name)
154 .where(DiscoveredRepository.actor == actor)
155 .where(DiscoveredRepository.discovered_at > cutoff)
156 ).all()
157 )
158
159 if not discovered:
160 discovered = set(self._discover_owned_repositories(actor))
161 if db is not None:
162 self._refresh_discovered_repositories(db, actor, discovered, discovered_at=now)
163
164 repositories = sorted(configured | tracked | discovered)
165 logger.info(
166 "git repositories selected for actor=%s count=%s configured=%s tracked=%s discovered=%s",
167 actor,
168 len(repositories),
169 len(configured),
170 len(tracked),
171 len(discovered),
172 )
173 return repositories
174
175 def fetch_backfill_batch(self, actor: str, cursor_state: dict | None = None) -> BackfillBatchResult:
176 return self._fetch_backfill_batch(actor, cursor_state, since=None)
177
178 def fetch_recent_backfill_batch(
179 self,
180 actor: str,
181 cursor_state: dict | None = None,
182 *,
183 since: datetime,
184 ) -> BackfillBatchResult:
185 return self._fetch_backfill_batch(actor, cursor_state, since=since)
186
187 def _fetch_backfill_batch(
188 self,
189 actor: str,
190 cursor_state: dict | None,
191 *,
192 since: datetime | None,
193 ) -> BackfillBatchResult:
194 state = {
195 "discovery_cursor": None,
196 "discovery_complete": False,
197 "repository_queue": sorted(
198 {
199 self._canonical_repository_name(actor, repository)
200 for repository in self.settings.git_tracked_repositories
201 }
202 ),
203 "current_repository": None,
204 }
205 if cursor_state:
206 state.update(copy.deepcopy(cursor_state))
207
208 if not state["discovery_complete"]:
209 data = self.client.execute(
210 USER_REPOSITORIES_QUERY,
211 {"username": actor.lstrip("~"), "cursor": state["discovery_cursor"]},
212 )
213 user = data.get("user") or {}
214 repositories_page = user.get("repositories") or {}
215 results = repositories_page.get("results") or []
216 state["discovery_cursor"] = repositories_page.get("cursor")
217 state["discovery_complete"] = not bool(state["discovery_cursor"])
218 known = set(state["repository_queue"])
219 current_repository = state.get("current_repository")
220 if current_repository:
221 known.add(current_repository["name"])
222 for repository in results:
223 if not isinstance(repository, dict):
224 continue
225 if repository.get("visibility") != "PUBLIC":
226 continue
227 name = repository.get("name")
228 repository_owner = ((repository.get("owner") or {}).get("canonicalName") or actor).strip()
229 if not name or not repository_owner:
230 continue
231 canonical_name = f"{repository_owner}/{name}"
232 if canonical_name not in known:
233 state["repository_queue"].append(canonical_name)
234 known.add(canonical_name)
235 state["repository_queue"] = sorted(state["repository_queue"])
236 logger.info(
237 "git backfill discovery actor=%s page_count=%s queue=%s next_cursor=%s",
238 actor,
239 len(results),
240 len(state["repository_queue"]),
241 bool(state["discovery_cursor"]),
242 )
243 complete = state["discovery_complete"] and not state["repository_queue"] and not state["current_repository"]
244 return BackfillBatchResult(events=[], cursor_state=state, complete=complete)
245
246 if state["current_repository"] is None:
247 if not state["repository_queue"]:
248 return BackfillBatchResult(events=[], cursor_state=None, complete=True)
249 state["current_repository"] = {
250 "name": state["repository_queue"].pop(0),
251 "reference_cursor": None,
252 "branches_loaded": False,
253 "branch_queue": [],
254 "current_branch": None,
255 }
256
257 repository_name = state["current_repository"]["name"]
258 owner, repo_name = self._split_repository(actor, repository_name)
259 current_repository = state["current_repository"]
260 current_repository.setdefault("reference_cursor", None)
261 current_repository.setdefault("branches_loaded", False)
262 current_repository.setdefault("branch_queue", [])
263 current_repository.setdefault("current_branch", None)
264 if "cursor" in current_repository:
265 current_repository.pop("cursor", None)
266
267 if not current_repository["branches_loaded"]:
268 data = self.client.execute(
269 REPOSITORY_BRANCHES_QUERY,
270 {"username": owner, "repoName": repo_name, "cursor": current_repository["reference_cursor"]},
271 )
272 user = data.get("user") or {}
273 repository = user.get("repository") or {}
274 references_page = repository.get("references") or {}
275 references = references_page.get("results") or []
276 current_repository["reference_cursor"] = references_page.get("cursor")
277 current_repository["branches_loaded"] = not bool(current_repository["reference_cursor"])
278 known_branches = set(current_repository["branch_queue"])
279 if current_repository["current_branch"]:
280 known_branches.add(current_repository["current_branch"]["name"])
281 for reference in references:
282 if not isinstance(reference, dict):
283 continue
284 branch_name = reference.get("name")
285 if not self._is_branch_reference(branch_name):
286 continue
287 if branch_name not in known_branches:
288 current_repository["branch_queue"].append(branch_name)
289 known_branches.add(branch_name)
290 current_repository["branch_queue"] = sorted(current_repository["branch_queue"])
291 logger.info(
292 "git backfill branch discovery actor=%s repository=%s page_count=%s branches=%s next_cursor=%s",
293 actor,
294 repository_name,
295 len(references),
296 len(current_repository["branch_queue"]),
297 bool(current_repository["reference_cursor"]),
298 )
299 if (
300 current_repository["branches_loaded"]
301 and not current_repository["branch_queue"]
302 and not current_repository["current_branch"]
303 ):
304 state["current_repository"] = None
305 complete = state["discovery_complete"] and not state["repository_queue"] and not state["current_repository"]
306 return BackfillBatchResult(events=[], cursor_state=None if complete else state, complete=complete)
307
308 if current_repository["current_branch"] is None:
309 if not current_repository["branch_queue"]:
310 state["current_repository"] = None
311 complete = state["discovery_complete"] and not state["repository_queue"]
312 return BackfillBatchResult(events=[], cursor_state=None if complete else state, complete=complete)
313 current_repository["current_branch"] = {"name": current_repository["branch_queue"].pop(0), "cursor": None}
314
315 current_branch = current_repository["current_branch"]
316 data = self.client.execute(
317 REPOSITORY_LOG_QUERY,
318 {
319 "username": owner,
320 "repoName": repo_name,
321 "cursor": current_branch["cursor"],
322 "from": current_branch["name"],
323 },
324 )
325 user = data.get("user") or {}
326 repository = user.get("repository") or {}
327 log_page = repository.get("log") or {}
328 commits = log_page.get("results") or []
329 next_cursor = log_page.get("cursor")
330 events: list[NormalizedEvent] = []
331 stop_repository = False
332 for commit in commits:
333 if not isinstance(commit, dict):
334 continue
335 commit_time = parse_datetime((commit.get("author") or {}).get("time"))
336 if since is not None and commit_time < since:
337 stop_repository = True
338 break
339 normalized = self._normalize_commit(actor=actor, repo_name=repo_name, commit=commit)
340 if normalized is not None:
341 events.append(normalized)
342 logger.info(
343 "git backfill actor=%s repository=%s branch=%s commits=%s next_cursor=%s",
344 actor,
345 repository_name,
346 current_branch["name"],
347 len(commits),
348 bool(next_cursor),
349 )
350 if next_cursor and not stop_repository:
351 current_branch["cursor"] = next_cursor
352 else:
353 current_repository["current_branch"] = None
354
355 complete = state["discovery_complete"] and not state["repository_queue"] and not state["current_repository"]
356 return BackfillBatchResult(events=events, cursor_state=None if complete else state, complete=complete)
357
358 def _discover_owned_repositories(self, actor: str) -> list[str]:
359 owner = actor.lstrip("~")
360 repositories: list[str] = []
361 cursor: str | None = None
362
363 for _ in range(50):
364 data = self.client.execute(
365 USER_REPOSITORIES_QUERY,
366 {"username": owner, "cursor": cursor},
367 )
368 user = data.get("user") or {}
369 repositories_page = user.get("repositories") or {}
370 results = repositories_page.get("results") or []
371 cursor = repositories_page.get("cursor")
372 logger.info(
373 "git repository discovery actor=%s page_count=%s next_cursor=%s",
374 actor,
375 len(results),
376 bool(cursor),
377 )
378
379 for repository in results:
380 if not isinstance(repository, dict):
381 continue
382 if repository.get("visibility") != "PUBLIC":
383 continue
384 name = repository.get("name")
385 repository_owner = ((repository.get("owner") or {}).get("canonicalName") or actor).strip()
386 if not name or not repository_owner:
387 continue
388 repo_name = f"{repository_owner}/{name}"
389 repositories.append(repo_name)
390
391 if not cursor:
392 break
393
394 return repositories
395
396 @staticmethod
397 def _refresh_discovered_repositories(
398 db: Session,
399 actor: str,
400 repositories: set[str],
401 *,
402 discovered_at: datetime,
403 ) -> None:
404 db.execute(delete(DiscoveredRepository).where(DiscoveredRepository.actor == actor))
405 for repository in sorted(repositories):
406 db.add(
407 DiscoveredRepository(
408 actor=actor,
409 name=repository,
410 discovered_at=discovered_at,
411 )
412 )
413 db.flush()
414
415 @staticmethod
416 def _canonical_repository_name(default_actor: str, repository: str) -> str:
417 owner, repo_name = GitIngestionService._split_repository(default_actor, repository)
418 return f"~{owner}/{repo_name}"
419
420 @staticmethod
421 def _split_repository(default_actor: str, repository: str) -> tuple[str, str]:
422 if "/" in repository:
423 owner, repo_name = repository.split("/", 1)
424 canonical_owner = owner if owner.startswith("~") else f"~{owner}"
425 return canonical_owner.lstrip("~"), repo_name
426 return default_actor.lstrip("~"), repository
427
428 def _fetch_repository_commits(
429 self,
430 *,
431 actor: str,
432 owner: str,
433 repo_name: str,
434 since: datetime,
435 ) -> list[NormalizedEvent]:
436 events: list[NormalizedEvent] = []
437 seen_event_uids: set[str] = set()
438 branches = self._fetch_repository_branches(owner=owner, repo_name=repo_name)
439 if not branches:
440 logger.info("git repository=%s/%s has no branch references", owner, repo_name)
441 return events
442
443 for branch in branches:
444 branch_events = self._fetch_repository_branch_commits(
445 actor=actor,
446 owner=owner,
447 repo_name=repo_name,
448 branch=branch,
449 since=since,
450 )
451 for event in branch_events:
452 if event.external_uid in seen_event_uids:
453 continue
454 seen_event_uids.add(event.external_uid)
455 events.append(event)
456
457 return events
458
459 def _fetch_repository_branches(self, *, owner: str, repo_name: str) -> list[str]:
460 branches: list[str] = []
461 cursor: str | None = None
462
463 for _ in range(50):
464 data = self.client.execute(
465 REPOSITORY_BRANCHES_QUERY,
466 {"username": owner, "repoName": repo_name, "cursor": cursor},
467 )
468 user = data.get("user") or {}
469 repository = user.get("repository") or {}
470 references_page = repository.get("references") or {}
471 references = references_page.get("results") or []
472 cursor = references_page.get("cursor")
473 logger.info(
474 "git repository=%s/%s branch page count=%s next_cursor=%s",
475 owner,
476 repo_name,
477 len(references),
478 bool(cursor),
479 )
480
481 for reference in references:
482 if not isinstance(reference, dict):
483 continue
484 branch_name = reference.get("name")
485 if self._is_branch_reference(branch_name):
486 branches.append(branch_name)
487
488 if not cursor:
489 break
490
491 return sorted(set(branches))
492
493 def _fetch_repository_branch_commits(
494 self,
495 *,
496 actor: str,
497 owner: str,
498 repo_name: str,
499 branch: str,
500 since: datetime,
501 ) -> list[NormalizedEvent]:
502 events: list[NormalizedEvent] = []
503 cursor: str | None = None
504
505 for _ in range(50):
506 data = self.client.execute(
507 REPOSITORY_LOG_QUERY,
508 {"username": owner, "repoName": repo_name, "cursor": cursor, "from": branch},
509 )
510 user = data.get("user") or {}
511 repository = user.get("repository") or {}
512 log_page = repository.get("log") or {}
513 commits = log_page.get("results") or []
514 cursor = log_page.get("cursor")
515 logger.info(
516 "git repository=%s/%s branch=%s commit page count=%s next_cursor=%s",
517 owner,
518 repo_name,
519 branch,
520 len(commits),
521 bool(cursor),
522 )
523
524 stop_paging = False
525 for commit in commits:
526 if not isinstance(commit, dict):
527 continue
528 commit_time = parse_datetime((commit.get("author") or {}).get("time"))
529 if commit_time < since:
530 stop_paging = True
531 logger.info(
532 "git commit %s skipped because commit_time=%s is before since=%s",
533 commit.get("shortId") or commit.get("id"),
534 commit_time.isoformat(),
535 since.isoformat(),
536 )
537 continue
538
539 normalized = self._normalize_commit(actor=actor, repo_name=repo_name, commit=commit)
540 if normalized is not None:
541 logger.info(
542 "git commit accepted repo=%s shortId=%s author=%s email=%s",
543 repo_name,
544 commit.get("shortId"),
545 (commit.get("author") or {}).get("name"),
546 (commit.get("author") or {}).get("email"),
547 )
548 events.append(normalized)
549 else:
550 logger.info(
551 "git commit skipped repo=%s shortId=%s author=%s email=%s",
552 repo_name,
553 commit.get("shortId"),
554 (commit.get("author") or {}).get("name"),
555 (commit.get("author") or {}).get("email"),
556 )
557
558 if stop_paging or not cursor:
559 break
560
561 return events
562
563 @staticmethod
564 def _is_branch_reference(reference_name: Any) -> bool:
565 return isinstance(reference_name, str) and reference_name.startswith("refs/heads/")
566
567 def _normalize_commit(
568 self,
569 *,
570 actor: str,
571 repo_name: str,
572 commit: dict[str, Any],
573 ) -> NormalizedEvent | None:
574 author = commit.get("author") or {}
575 candidate_aliases = [
576 actor,
577 author.get("email", ""),
578 author.get("name", ""),
579 ]
580 matched_actor = None
581 for candidate in candidate_aliases:
582 canonical = self.identity_resolver.canonicalize(candidate)
583 if canonical == actor:
584 matched_actor = canonical
585 break
586
587 if matched_actor is None:
588 return None
589
590 commit_id = str(commit["id"])
591 commit_time = parse_datetime(author["time"])
592 return NormalizedEvent(
593 service=self.service_name,
594 event_type="commit",
595 actor=matched_actor,
596 repo_name=repo_name,
597 resource_id=commit_id,
598 external_uid=f"git:{repo_name}:{commit_id}",
599 occurred_at=commit_time,
600 weight=self.settings.event_weights["commit"],
601 raw_payload_json=commit,
602 )