| @@ -21,6 +21,23 @@ from srht_contrib.utils.identity import ActorIdentityResolver |
| 21 | 21 | logger = logging.getLogger(__name__) |
| 22 | 22 | |
| 23 | 23 | |
| 24 | REPOSITORY_BRANCHES_QUERY = """ |
| 25 | query 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 | |
| 24 | 41 | REPOSITORY_LOG_QUERY = """ |
| 25 | 42 | query RepositoryLog($username: String!, $repoName: String!, $cursor: Cursor, $from: String) { |
| 26 | 43 | user(username: $username) { |
| @@ -29,10 +46,6 @@ query RepositoryLog($username: String!, $repoName: String!, $cursor: Cursor, $fr |
| 29 | 46 | owner { |
| 30 | 47 | canonicalName |
| 31 | 48 | } |
| 32 | | HEAD { |
| 33 | | name |
| 34 | | target |
| 35 | | } |
| 36 | 49 | log(cursor: $cursor, from: $from) { |
| 37 | 50 | results { |
| 38 | 51 | id |
| @@ -233,17 +246,80 @@ class GitIngestionService: |
| 233 | 246 | if state["current_repository"] is None: |
| 234 | 247 | if not state["repository_queue"]: |
| 235 | 248 | return BackfillBatchResult(events=[], cursor_state=None, complete=True) |
| 236 | | state["current_repository"] = {"name": state["repository_queue"].pop(0), "cursor": None} |
| 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 | } |
| 237 | 256 | |
| 238 | 257 | repository_name = state["current_repository"]["name"] |
| 239 | 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"] |
| 240 | 316 | data = self.client.execute( |
| 241 | 317 | REPOSITORY_LOG_QUERY, |
| 242 | 318 | { |
| 243 | 319 | "username": owner, |
| 244 | 320 | "repoName": repo_name, |
| 245 | | "cursor": state["current_repository"]["cursor"], |
| 246 | | "from": "HEAD", |
| 321 | "cursor": current_branch["cursor"], |
| 322 | "from": current_branch["name"], |
| 247 | 323 | }, |
| 248 | 324 | ) |
| 249 | 325 | user = data.get("user") or {} |
| @@ -264,16 +340,17 @@ class GitIngestionService: |
| 264 | 340 | if normalized is not None: |
| 265 | 341 | events.append(normalized) |
| 266 | 342 | logger.info( |
| 267 | | "git backfill actor=%s repository=%s commits=%s next_cursor=%s", |
| 343 | "git backfill actor=%s repository=%s branch=%s commits=%s next_cursor=%s", |
| 268 | 344 | actor, |
| 269 | 345 | repository_name, |
| 346 | current_branch["name"], |
| 270 | 347 | len(commits), |
| 271 | 348 | bool(next_cursor), |
| 272 | 349 | ) |
| 273 | 350 | if next_cursor and not stop_repository: |
| 274 | | state["current_repository"]["cursor"] = next_cursor |
| 351 | current_branch["cursor"] = next_cursor |
| 275 | 352 | else: |
| 276 | | state["current_repository"] = None |
| 353 | current_repository["current_branch"] = None |
| 277 | 354 | |
| 278 | 355 | complete = state["discovery_complete"] and not state["repository_queue"] and not state["current_repository"] |
| 279 | 356 | return BackfillBatchResult(events=events, cursor_state=None if complete else state, complete=complete) |
| @@ -355,6 +432,72 @@ class GitIngestionService: |
| 355 | 432 | owner: str, |
| 356 | 433 | repo_name: str, |
| 357 | 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, |
| 358 | 501 | ) -> list[NormalizedEvent]: |
| 359 | 502 | events: list[NormalizedEvent] = [] |
| 360 | 503 | cursor: str | None = None |
| @@ -362,7 +505,7 @@ class GitIngestionService: |
| 362 | 505 | for _ in range(50): |
| 363 | 506 | data = self.client.execute( |
| 364 | 507 | REPOSITORY_LOG_QUERY, |
| 365 | | {"username": owner, "repoName": repo_name, "cursor": cursor, "from": "HEAD"}, |
| 508 | {"username": owner, "repoName": repo_name, "cursor": cursor, "from": branch}, |
| 366 | 509 | ) |
| 367 | 510 | user = data.get("user") or {} |
| 368 | 511 | repository = user.get("repository") or {} |
| @@ -370,9 +513,10 @@ class GitIngestionService: |
| 370 | 513 | commits = log_page.get("results") or [] |
| 371 | 514 | cursor = log_page.get("cursor") |
| 372 | 515 | logger.info( |
| 373 | | "git repository=%s/%s commit page count=%s next_cursor=%s", |
| 516 | "git repository=%s/%s branch=%s commit page count=%s next_cursor=%s", |
| 374 | 517 | owner, |
| 375 | 518 | repo_name, |
| 519 | branch, |
| 376 | 520 | len(commits), |
| 377 | 521 | bool(cursor), |
| 378 | 522 | ) |
| @@ -416,6 +560,10 @@ class GitIngestionService: |
| 416 | 560 | |
| 417 | 561 | return events |
| 418 | 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 | |
| 419 | 567 | def _normalize_commit( |
| 420 | 568 | self, |
| 421 | 569 | *, |