diff options
| author | Christian Cleberg <[email protected]> | 2026-05-06 22:09:56 -0500 |
|---|---|---|
| committer | Christian Cleberg <[email protected]> | 2026-05-06 22:09:56 -0500 |
| commit | 1ae229d64bee19e94bbfd1872879b86496d68104 (patch) | |
| tree | cb187f435df234f0bc2f77ebe1ddae5842c6d0ac | |
| parent | fdfb7831fe7dba6d17a1697c7cbb731fd008ba7b (diff) | |
| download | hutch-stats-1ae229d64bee19e94bbfd1872879b86496d68104.tar.gz hutch-stats-1ae229d64bee19e94bbfd1872879b86496d68104.tar.bz2 hutch-stats-1ae229d64bee19e94bbfd1872879b86496d68104.zip | |
feat: index git activity across all SourceHut branches
| -rw-r--r-- | src/srht_contrib/services/git.py | 172 | ||||
| -rw-r--r-- | tests/test_git_repository_cache.py | 14 | ||||
| -rw-r--r-- | tests/test_ingestion.py | 98 |
3 files changed, 235 insertions, 49 deletions
diff --git a/src/srht_contrib/services/git.py b/src/srht_contrib/services/git.py index ee5c05c..7b5fdf3 100644 --- a/src/srht_contrib/services/git.py +++ b/src/srht_contrib/services/git.py @@ -21,6 +21,23 @@ from srht_contrib.utils.identity import ActorIdentityResolver logger = logging.getLogger(__name__) +REPOSITORY_BRANCHES_QUERY = """ +query RepositoryBranches($username: String!, $repoName: String!, $cursor: Cursor) { + user(username: $username) { + repository(name: $repoName) { + references(cursor: $cursor) { + results { + name + target + } + cursor + } + } + } +} +""".strip() + + REPOSITORY_LOG_QUERY = """ query RepositoryLog($username: String!, $repoName: String!, $cursor: Cursor, $from: String) { user(username: $username) { @@ -29,10 +46,6 @@ query RepositoryLog($username: String!, $repoName: String!, $cursor: Cursor, $fr owner { canonicalName } - HEAD { - name - target - } log(cursor: $cursor, from: $from) { results { id @@ -233,17 +246,80 @@ class GitIngestionService: if state["current_repository"] is None: if not state["repository_queue"]: return BackfillBatchResult(events=[], cursor_state=None, complete=True) - state["current_repository"] = {"name": state["repository_queue"].pop(0), "cursor": None} + state["current_repository"] = { + "name": state["repository_queue"].pop(0), + "reference_cursor": None, + "branches_loaded": False, + "branch_queue": [], + "current_branch": None, + } repository_name = state["current_repository"]["name"] owner, repo_name = self._split_repository(actor, repository_name) + current_repository = state["current_repository"] + current_repository.setdefault("reference_cursor", None) + current_repository.setdefault("branches_loaded", False) + current_repository.setdefault("branch_queue", []) + current_repository.setdefault("current_branch", None) + if "cursor" in current_repository: + current_repository.pop("cursor", None) + + if not current_repository["branches_loaded"]: + data = self.client.execute( + REPOSITORY_BRANCHES_QUERY, + {"username": owner, "repoName": repo_name, "cursor": current_repository["reference_cursor"]}, + ) + user = data.get("user") or {} + repository = user.get("repository") or {} + references_page = repository.get("references") or {} + references = references_page.get("results") or [] + current_repository["reference_cursor"] = references_page.get("cursor") + current_repository["branches_loaded"] = not bool(current_repository["reference_cursor"]) + known_branches = set(current_repository["branch_queue"]) + if current_repository["current_branch"]: + known_branches.add(current_repository["current_branch"]["name"]) + for reference in references: + if not isinstance(reference, dict): + continue + branch_name = reference.get("name") + if not self._is_branch_reference(branch_name): + continue + if branch_name not in known_branches: + current_repository["branch_queue"].append(branch_name) + known_branches.add(branch_name) + current_repository["branch_queue"] = sorted(current_repository["branch_queue"]) + logger.info( + "git backfill branch discovery actor=%s repository=%s page_count=%s branches=%s next_cursor=%s", + actor, + repository_name, + len(references), + len(current_repository["branch_queue"]), + bool(current_repository["reference_cursor"]), + ) + if ( + current_repository["branches_loaded"] + and not current_repository["branch_queue"] + and not current_repository["current_branch"] + ): + state["current_repository"] = None + complete = state["discovery_complete"] and not state["repository_queue"] and not state["current_repository"] + return BackfillBatchResult(events=[], cursor_state=None if complete else state, complete=complete) + + if current_repository["current_branch"] is None: + if not current_repository["branch_queue"]: + state["current_repository"] = None + complete = state["discovery_complete"] and not state["repository_queue"] + return BackfillBatchResult(events=[], cursor_state=None if complete else state, complete=complete) + current_repository["current_branch"] = {"name": current_repository["branch_queue"].pop(0), "cursor": None} + + current_branch = current_repository["current_branch"] data = self.client.execute( REPOSITORY_LOG_QUERY, { "username": owner, "repoName": repo_name, - "cursor": state["current_repository"]["cursor"], - "from": "HEAD", + "cursor": current_branch["cursor"], + "from": current_branch["name"], }, ) user = data.get("user") or {} @@ -264,16 +340,17 @@ class GitIngestionService: if normalized is not None: events.append(normalized) logger.info( - "git backfill actor=%s repository=%s commits=%s next_cursor=%s", + "git backfill actor=%s repository=%s branch=%s commits=%s next_cursor=%s", actor, repository_name, + current_branch["name"], len(commits), bool(next_cursor), ) if next_cursor and not stop_repository: - state["current_repository"]["cursor"] = next_cursor + current_branch["cursor"] = next_cursor else: - state["current_repository"] = None + current_repository["current_branch"] = None complete = state["discovery_complete"] and not state["repository_queue"] and not state["current_repository"] return BackfillBatchResult(events=events, cursor_state=None if complete else state, complete=complete) @@ -357,12 +434,78 @@ class GitIngestionService: since: datetime, ) -> list[NormalizedEvent]: events: list[NormalizedEvent] = [] + seen_event_uids: set[str] = set() + branches = self._fetch_repository_branches(owner=owner, repo_name=repo_name) + if not branches: + logger.info("git repository=%s/%s has no branch references", owner, repo_name) + return events + + for branch in branches: + branch_events = self._fetch_repository_branch_commits( + actor=actor, + owner=owner, + repo_name=repo_name, + branch=branch, + since=since, + ) + for event in branch_events: + if event.external_uid in seen_event_uids: + continue + seen_event_uids.add(event.external_uid) + events.append(event) + + return events + + def _fetch_repository_branches(self, *, owner: str, repo_name: str) -> list[str]: + branches: list[str] = [] + cursor: str | None = None + + for _ in range(50): + data = self.client.execute( + REPOSITORY_BRANCHES_QUERY, + {"username": owner, "repoName": repo_name, "cursor": cursor}, + ) + user = data.get("user") or {} + repository = user.get("repository") or {} + references_page = repository.get("references") or {} + references = references_page.get("results") or [] + cursor = references_page.get("cursor") + logger.info( + "git repository=%s/%s branch page count=%s next_cursor=%s", + owner, + repo_name, + len(references), + bool(cursor), + ) + + for reference in references: + if not isinstance(reference, dict): + continue + branch_name = reference.get("name") + if self._is_branch_reference(branch_name): + branches.append(branch_name) + + if not cursor: + break + + return sorted(set(branches)) + + def _fetch_repository_branch_commits( + self, + *, + actor: str, + owner: str, + repo_name: str, + branch: 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, "from": "HEAD"}, + {"username": owner, "repoName": repo_name, "cursor": cursor, "from": branch}, ) user = data.get("user") or {} repository = user.get("repository") or {} @@ -370,9 +513,10 @@ class GitIngestionService: commits = log_page.get("results") or [] cursor = log_page.get("cursor") logger.info( - "git repository=%s/%s commit page count=%s next_cursor=%s", + "git repository=%s/%s branch=%s commit page count=%s next_cursor=%s", owner, repo_name, + branch, len(commits), bool(cursor), ) @@ -416,6 +560,10 @@ class GitIngestionService: return events + @staticmethod + def _is_branch_reference(reference_name: Any) -> bool: + return isinstance(reference_name, str) and reference_name.startswith("refs/heads/") + def _normalize_commit( self, *, diff --git a/tests/test_git_repository_cache.py b/tests/test_git_repository_cache.py index 9326e6f..508f668 100644 --- a/tests/test_git_repository_cache.py +++ b/tests/test_git_repository_cache.py @@ -40,6 +40,19 @@ def make_settings(**overrides) -> Settings: return Settings(**values) +def branch_payload(*branches: str) -> dict: + return { + "user": { + "repository": { + "references": { + "results": [{"name": branch, "target": "abc123"} for branch in branches], + "cursor": None, + } + } + } + } + + def test_git_poll_reuses_cached_discovered_repositories(db_session) -> None: settings = make_settings() client = StubClient( @@ -58,6 +71,7 @@ def test_git_poll_reuses_cached_discovered_repositories(db_session) -> None: } } }, + "query RepositoryBranches": branch_payload("refs/heads/main"), "query RepositoryLog": { "user": { "repository": { diff --git a/tests/test_ingestion.py b/tests/test_ingestion.py index dc96354..87b956b 100644 --- a/tests/test_ingestion.py +++ b/tests/test_ingestion.py @@ -202,6 +202,19 @@ def make_settings(**overrides) -> Settings: return Settings(**values) +def branch_payload(*branches: str) -> dict: + return { + "user": { + "repository": { + "references": { + "results": [{"name": branch, "target": "abc123"} for branch in branches], + "cursor": None, + } + } + } + } + + def test_todo_ingestion_is_idempotent(db_session) -> None: settings = make_settings() payload = { @@ -431,7 +444,15 @@ def test_git_ingestion_normalizes_commit_aliases_and_repository_names(db_session StubClient(payload={"me": {"canonicalName": "~ccleberg"}, "events": {"results": [], "cursor": None}}), settings, ) - git_service = GitIngestionService(StubClient(payloads_by_query={"query RepositoryLog": git_payload}), settings) + git_service = GitIngestionService( + StubClient( + payloads_by_query={ + "query RepositoryBranches": branch_payload("refs/heads/main"), + "query RepositoryLog": git_payload, + } + ), + settings, + ) poller = PollerService(todo_service=todo_service, git_service=git_service, settings=settings) inserted = poller.poll_all(db_session, "~ccleberg") @@ -463,6 +484,7 @@ def test_git_ingestion_auto_discovers_owned_repositories(db_session) -> None: } } }, + "query RepositoryBranches": branch_payload("refs/heads/main"), "query RepositoryLog": { "user": { "repository": { @@ -507,53 +529,55 @@ def test_git_ingestion_auto_discovers_owned_repositories(db_session) -> None: assert any("query UserRepositories" in call[0] for call in client.calls) -def test_git_ingestion_reads_repository_log_from_default_head_branch(db_session) -> None: +def test_git_ingestion_reads_repository_logs_from_all_branches(db_session) -> None: settings = make_settings( ACTOR_ALIASES_JSON={"~ccleberg": ["[email protected]", "Chris Cleberg"]}, GIT_TRACKED_REPOSITORIES=["Hutch"], ) - git_payload = { - "user": { - "repository": { - "name": "Hutch", - "owner": {"canonicalName": "~ccleberg"}, - "HEAD": {"name": "refs/heads/trunk", "target": "abc123"}, - "log": { - "results": [ - { - "id": "abc123", - "shortId": "abc123", - "author": { - "name": "Chris Cleberg", - "email": "[email protected]", - "time": "2026-03-30T12:00:00Z", - }, - "committer": { - "name": "Chris Cleberg", - "email": "[email protected]", - "time": "2026-03-30T12:00:00Z", - }, - "message": "Commit on non-standard default branch", - } - ], - "cursor": None, - }, - } + client = StubClient( + payloads_by_query={ + "query RepositoryBranches": branch_payload( + "refs/heads/main", + "refs/heads/trunk", + "refs/tags/v1.0.0", + ), + "query RepositoryLog": { + "user": { + "repository": { + "name": "Hutch", + "owner": {"canonicalName": "~ccleberg"}, + "log": { + "results": [ + { + "id": "abc123", + "shortId": "abc123", + "author": { + "name": "Chris Cleberg", + "email": "[email protected]", + "time": "2026-03-30T12:00:00Z", + }, + "committer": { + "name": "Chris Cleberg", + "email": "[email protected]", + "time": "2026-03-30T12:00:00Z", + }, + "message": "Commit reachable from more than one branch", + } + ], + "cursor": None, + }, + } + } + }, } - } - client = StubClient(payloads_by_query={"query RepositoryLog": git_payload}) + ) git_service = GitIngestionService(client, settings) result = git_service.fetch_recent_events("~ccleberg", since=datetime(2026, 3, 1, tzinfo=UTC)) repository_log_calls = [call for call in client.calls if "query RepositoryLog" in call[0]] assert len(result.events) == 1 - assert repository_log_calls[0][1] == { - "username": "ccleberg", - "repoName": "Hutch", - "cursor": None, - "from": "HEAD", - } + assert [call[1]["from"] for call in repository_log_calls] == ["refs/heads/main", "refs/heads/trunk"] def test_sync_overlap_reuses_cursor_window_and_suppresses_duplicates(db_session) -> None: |
