Skip to content

Commit f37f071

Browse files
authored
feat: support edge case of forked repo having frequent force pushes (CM-745) (#3663)
1 parent 84f6b15 commit f37f071

4 files changed

Lines changed: 39 additions & 52 deletions

File tree

services/apps/git_integration/src/crowdgit/database/crud.py

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@ async def insert_repository(url: str, priority: int = 0) -> str:
3232
async def get_repository_by_url(url: str) -> dict[str, Any] | None:
3333
"""Get repository by URL"""
3434
sql_query = """
35-
SELECT id, url, state, priority, "lastProcessedAt", "lockedAt", "createdAt", "updatedAt", "maintainerFile", "forkedFrom"
35+
SELECT id, url, state, priority, "lastProcessedAt", "lockedAt", "createdAt", "updatedAt", "maintainerFile", "forkedFrom", "stuckRequiresReOnboard"
3636
FROM git.repositories
3737
WHERE url = $1 AND "deletedAt" IS NULL
3838
"""
@@ -49,7 +49,7 @@ async def get_recently_processed_repository_by_url(url: str) -> Repository | Non
4949
Used to check if a repository needs reprocessing based on the update interval.
5050
"""
5151
sql_query = """
52-
SELECT id, url, state, priority, "lastProcessedAt", "lockedAt", "createdAt", "updatedAt", "maintainerFile", "forkedFrom", "segmentId"
52+
SELECT id, url, state, priority, "lastProcessedAt", "lockedAt", "createdAt", "updatedAt", "maintainerFile", "forkedFrom", "segmentId", "stuckRequiresReOnboard"
5353
FROM git.repositories
5454
WHERE url = $1
5555
AND "deletedAt" IS NULL
@@ -88,7 +88,7 @@ async def acquire_onboarding_repo() -> Repository | None:
8888
LIMIT 1
8989
FOR UPDATE SKIP LOCKED
9090
)
91-
RETURNING id, url, state, priority, "lastProcessedAt", "lastProcessedCommit", "lockedAt", "createdAt", "updatedAt", "segmentId", "integrationId", "maintainerFile", "lastMaintainerRunAt", "branch", "forkedFrom"
91+
RETURNING id, url, state, priority, "lastProcessedAt", "lastProcessedCommit", "lockedAt", "createdAt", "updatedAt", "segmentId", "integrationId", "maintainerFile", "lastMaintainerRunAt", "branch", "forkedFrom", "stuckRequiresReOnboard"
9292
"""
9393
return await acquire_repository(
9494
onboarding_repo_sql_query,
@@ -138,7 +138,7 @@ async def acquire_recurrent_repo() -> Repository | None:
138138
LIMIT 1
139139
FOR UPDATE SKIP LOCKED
140140
)
141-
RETURNING id, url, state, priority, "lastProcessedAt", "lastProcessedCommit", "lockedAt", "createdAt", "updatedAt", "segmentId", "integrationId", "maintainerFile", "lastMaintainerRunAt", "branch", "forkedFrom"
141+
RETURNING id, url, state, priority, "lastProcessedAt", "lastProcessedCommit", "lockedAt", "createdAt", "updatedAt", "segmentId", "integrationId", "maintainerFile", "lastMaintainerRunAt", "branch", "forkedFrom", "stuckRequiresReOnboard"
142142
"""
143143
states_to_exclude = (
144144
RepositoryState.PENDING,

services/apps/git_integration/src/crowdgit/models/repository.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,10 @@ class Repository(BaseModel):
4141
parent_repo: Repository | None = Field(
4242
None, description="The parent repository (in case of fork) object from our database"
4343
)
44+
stuck_requires_re_onboard: bool = Field(
45+
default=False,
46+
description="Indicates if the stuck repository is resolved by a re-onboarding",
47+
)
4448
created_at: datetime = Field(..., description="Creation timestamp")
4549
updated_at: datetime = Field(..., description="Last update timestamp")
4650

@@ -67,6 +71,7 @@ def from_db(cls, db_data: dict[str, Any]) -> Repository:
6771
"maintainerFile": "maintainer_file",
6872
"lastMaintainerRunAt": "last_maintainer_run_at",
6973
"forkedFrom": "forked_from",
74+
"stuckRequiresReOnboard": "stuck_requires_re_onboard",
7075
}
7176
for db_field, model_field in field_mapping.items():
7277
if db_field in repo_data:

services/apps/git_integration/src/crowdgit/services/commit/commit_service.py

Lines changed: 29 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -137,15 +137,7 @@ async def process_single_batch_commits(
137137
repository.last_processed_commit,
138138
)
139139

140-
await self._process_activities_from_commits(
141-
raw_commits,
142-
batch_info.repo_path,
143-
batch_info.edge_commit,
144-
batch_info.remote,
145-
repository.segment_id,
146-
repository.integration_id,
147-
repository.parent_repo,
148-
)
140+
await self._process_activities_from_commits(raw_commits, batch_info, repository)
149141

150142
batch_end_time = time.time()
151143
batch_time = round(batch_end_time - batch_start_time, 2)
@@ -616,14 +608,14 @@ def create_activities_from_commit(
616608

617609
return activities_db, activities_queue
618610

619-
async def _filter_parent_repo_activities(
611+
async def _filter_existing_activities(
620612
self,
621613
activities_db: list[tuple],
622614
activities_queue: list[dict],
623-
parent_repo: Repository,
615+
source_repo: Repository,
624616
) -> tuple[list[tuple], list[dict], int]:
625617
"""
626-
Filter out activities that exist in parent repo for forked repositories.
618+
Filter out activities that exist in specific repo, used for both forked and frequently reonboarded repos.
627619
Done in post-processing phase using batch lookup to avoid N+1 queries.
628620
629621
Returns: (filtered_activities_db, filtered_activities_queue, skipped_activities_count)
@@ -639,8 +631,8 @@ async def _filter_parent_repo_activities(
639631
# Batch check which activities exist in parent repo
640632
parent_source_ids = await batch_check_parent_activities(
641633
activity_keys,
642-
parent_repo.url,
643-
parent_repo.segment_id,
634+
source_repo.url,
635+
source_repo.segment_id,
644636
)
645637

646638
if not parent_source_ids:
@@ -664,20 +656,16 @@ async def _filter_parent_repo_activities(
664656

665657
if skipped_activities_count > 0:
666658
self.logger.info(
667-
f"Filtered out {skipped_activities_count} activities from parent repo {parent_repo.url}"
659+
f"Filtered out {skipped_activities_count} existing activity from {source_repo.url}"
668660
)
669661

670662
return filtered_activities_db, filtered_activities_queue, skipped_activities_count
671663

672664
async def process_commits_chunk(
673665
self,
674666
commit_texts_chunk: list[str | None],
675-
repo_path: str,
676-
edge_commit_hash: str | None,
677-
remote: str,
678-
segment_id: str,
679-
integration_id: str,
680-
parent_repo: Repository | None,
667+
batch_info: CloneBatchInfo,
668+
repository: Repository,
681669
) -> None:
682670
"""
683671
Process a chunk of raw commit texts into activities and write them to DB and Kafka.
@@ -697,15 +685,15 @@ async def process_commits_chunk(
697685
commit = None
698686

699687
for full_commit_text in commit_texts_chunk:
700-
if self.should_skip_commit(full_commit_text, edge_commit_hash):
688+
if self.should_skip_commit(full_commit_text, batch_info.edge_commit):
701689
continue
702690
commit_text, numstats_text = full_commit_text.split(self.NUMSTAT_SPLITTER)
703691
commit_lines = commit_text.strip().splitlines()
704692
del full_commit_text
705693
del commit_text
706694
if not self._validate_commit_structure(commit_lines):
707695
self.logger.warning(
708-
f"Invalid commit structure in {repo_path}: {len(commit_lines)} fields"
696+
f"Invalid commit structure in {batch_info.repo_path}: {len(commit_lines)} fields"
709697
)
710698
bad_commits += 1
711699
del commit_lines
@@ -716,7 +704,7 @@ async def process_commits_chunk(
716704
commit = self._construct_commit_dict(commit_lines, numstats_text)
717705
if self._validate_commit_data(commit):
718706
activity_db_records, activity_kafka = self.create_activities_from_commit(
719-
remote, commit, segment_id, integration_id
707+
batch_info.remote, commit, repository.segment_id, repository.integration_id
720708
)
721709
activities_db.extend(activity_db_records)
722710
activities_queue.extend(activity_kafka)
@@ -727,7 +715,7 @@ async def process_commits_chunk(
727715
bad_commits += 1
728716

729717
except Exception as e:
730-
self.logger.warning(f"Failed to parse commit in {repo_path}: {e}")
718+
self.logger.warning(f"Failed to parse commit in {batch_info.repo_path}: {e}")
731719
bad_commits += 1
732720
continue
733721
finally:
@@ -737,17 +725,26 @@ async def process_commits_chunk(
737725

738726
# Filter out activities from parent repo (for forks)
739727
skipped_activities = 0
740-
if parent_repo:
728+
if repository.parent_repo:
741729
(
742730
activities_db,
743731
activities_queue,
744732
skipped_activities,
745-
) = await self._filter_parent_repo_activities(
746-
activities_db, activities_queue, parent_repo
733+
) = await self._filter_existing_activities(
734+
activities_db, activities_queue, repository.parent_repo
735+
)
736+
if repository.stuck_requires_re_onboard:
737+
self.logger.info(
738+
f"Frequent re-onboardings detected! excluding existing activities from repo: {repository.url}"
747739
)
740+
(
741+
activities_db,
742+
activities_queue,
743+
skipped_activities,
744+
) = await self._filter_existing_activities(activities_db, activities_queue, repository)
748745

749746
self.logger.info(
750-
f"Processed {processed_commits} commits, skipped {bad_commits} invalid commits, filtered {skipped_activities} activities from parent repo in {repo_path}"
747+
f"Processed {processed_commits} commits, skipped {bad_commits} invalid commits, filtered {skipped_activities} activities from parent repo in {batch_info.repo_path}"
751748
)
752749
# Update metrics context
753750
if self._metrics_context:
@@ -766,14 +763,7 @@ async def process_commits_chunk(
766763
del activities_db, activities_queue
767764

768765
async def _process_activities_from_commits(
769-
self,
770-
raw_commits: str,
771-
repo_path: str,
772-
edge_commit_hash: str | None,
773-
remote: str,
774-
segment_id: str,
775-
integration_id: str,
776-
parent_repo: Repository | None = None,
766+
self, raw_commits: str, batch_info: CloneBatchInfo, repository: Repository
777767
):
778768
"""
779769
Parse raw git log output, process commits into activities, and save to database.
@@ -815,12 +805,8 @@ async def process_single_chunk(chunk_start_idx: int, chunk_end_idx: int):
815805
# Process chunk and write to DB/Kafka
816806
await self.process_commits_chunk(
817807
chunk,
818-
repo_path,
819-
edge_commit_hash,
820-
remote,
821-
segment_id,
822-
integration_id,
823-
parent_repo,
808+
batch_info,
809+
repository,
824810
)
825811
completed_chunks += 1
826812
self.logger.info(f"Progress: {completed_chunks}/{total_chunks} chunks")

services/apps/git_integration/src/crowdgit/worker/repository_worker.py

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -115,7 +115,7 @@ async def _ensure_repo_not_stuck(self, repository: Repository):
115115
logger.warning(
116116
f"Repo {repository.url} is stuck for {processing_duration_hours} hours!"
117117
)
118-
if repository.forked_from == repository.url:
118+
if repository.stuck_requires_re_onboard:
119119
logger.warning(
120120
f"Repo {repository.url} is stuck due to force-push or dangling commit. Will be re-onboarded"
121121
)
@@ -197,10 +197,6 @@ async def _validate_and_get_parent_repo(self, repository: Repository) -> Reposit
197197
if not repository.forked_from:
198198
return None
199199

200-
if repository.forked_from == repository.url:
201-
# EDGE CASE: not a fork but repo get reonboarded a lot and we treat it as a "fork" to avoid producing tons of duplicate activities
202-
return repository.forked_from
203-
204200
logger.info(
205201
f"Repository {repository.url} is forked from {repository.forked_from}, validating parent repo..."
206202
)

0 commit comments

Comments
 (0)