Skip to content

Commit 17a99a3

Browse files
committed
fix: resolve pr review comments
Signed-off-by: Yeganathan S <63534555+skwowet@users.noreply.github.com>
1 parent 3e1c41c commit 17a99a3

3 files changed

Lines changed: 129 additions & 98 deletions

File tree

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

Lines changed: 17 additions & 47 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,11 @@
11
from datetime import datetime, timezone
22

33
from loguru import logger
4-
from pydantic import TypeAdapter, ValidationError
54
from tenacity import retry, retry_if_exception_type, stop_after_attempt, wait_fixed
65

76
from crowdgit.enums import RepositoryPriority, RepositoryState
87
from crowdgit.errors import RepoLockingError
9-
from crowdgit.models.affiliation_info import AffiliationInfoItem
8+
from crowdgit.models.affiliation_info import RepoAffiliationRegistry
109
from crowdgit.models.repository import Repository
1110
from crowdgit.models.service_execution import ServiceExecution
1211
from crowdgit.settings import (
@@ -528,24 +527,7 @@ async def save_service_execution(service_execution: ServiceExecution) -> None:
528527
# Do not re-raise - we don't want metrics saving to disrupt main operations
529528

530529

531-
_AFFILIATION_SNAPSHOT_ADAPTER = TypeAdapter(list[AffiliationInfoItem])
532-
533-
534-
def parse_affiliation_snapshot(snapshot) -> list[AffiliationInfoItem] | None:
535-
if isinstance(snapshot, dict) and "affiliations" in snapshot:
536-
snapshot = snapshot["affiliations"]
537-
try:
538-
return _AFFILIATION_SNAPSHOT_ADAPTER.validate_python(snapshot)
539-
except ValidationError as error:
540-
logger.warning(f"Invalid affiliation snapshot in registry, will re-parse: {error}")
541-
return None
542-
543-
544-
def dump_affiliation_snapshot(affiliations: list[AffiliationInfoItem]) -> list[dict]:
545-
return [item.model_dump() for item in affiliations]
546-
547-
548-
async def get_repo_affiliation_registry(repo_id: str):
530+
async def get_repo_affiliation_registry(repo_id: str) -> RepoAffiliationRegistry | None:
549531
sql_query = """
550532
SELECT "filePath", "fileHash", "status", "snapshot", "lastRunAt"
551533
FROM git."repoAffiliationRegistry"
@@ -556,28 +538,12 @@ async def get_repo_affiliation_registry(repo_id: str):
556538
return None
557539

558540
row = dict(result)
559-
snapshot = row.get("snapshot")
560-
if snapshot is not None:
561-
snapshot = parse_affiliation_snapshot(snapshot)
562-
563-
return {
564-
"file_path": row.get("filePath"),
565-
"file_hash": row.get("fileHash"),
566-
"status": row.get("status"),
567-
"snapshot": snapshot,
568-
"last_run_at": row.get("lastRunAt"),
569-
}
541+
row["repoId"] = repo_id
542+
return RepoAffiliationRegistry.from_db(row)
570543

571544

572-
async def upsert_repo_affiliation_registry(
573-
repo_id: str,
574-
*,
575-
file_path: str | None,
576-
file_hash: str | None,
577-
status: str,
578-
snapshot: list[AffiliationInfoItem] | None,
579-
) -> None:
580-
snapshot_json = dump_affiliation_snapshot(snapshot) if snapshot is not None else None
545+
async def upsert_repo_affiliation_registry(registry: RepoAffiliationRegistry) -> None:
546+
snapshot_json = registry.snapshot_for_db()
581547
sql_query = """
582548
INSERT INTO git."repoAffiliationRegistry" (
583549
"repoId", "filePath", "fileHash", "status", "snapshot", "lastRunAt", "updatedAt"
@@ -593,7 +559,13 @@ async def upsert_repo_affiliation_registry(
593559
"""
594560
await execute(
595561
sql_query,
596-
(repo_id, file_path, file_hash, status, snapshot_json),
562+
(
563+
registry.repo_id,
564+
registry.file_path,
565+
registry.file_hash,
566+
registry.status,
567+
snapshot_json,
568+
),
597569
)
598570

599571

@@ -750,9 +722,9 @@ async def fetch_segment_affiliations(member_ids: list[str], segment_id: str) ->
750722
)
751723

752724

753-
async def insert_member_organizations(rows: list[dict]) -> int:
725+
async def insert_member_organizations(rows: list[dict]) -> None:
754726
if not rows:
755-
return 0
727+
return
756728

757729
sql_query = """
758730
INSERT INTO "memberOrganizations"(
@@ -780,12 +752,11 @@ async def insert_member_organizations(rows: list[dict]) -> int:
780752
for row in rows
781753
],
782754
)
783-
return len(rows)
784755

785756

786-
async def insert_member_segment_affiliations(rows: list[dict]) -> int:
757+
async def insert_member_segment_affiliations(rows: list[dict]) -> None:
787758
if not rows:
788-
return 0
759+
return
789760

790761
sql_query = """
791762
INSERT INTO "memberSegmentAffiliations"(
@@ -811,4 +782,3 @@ async def insert_member_segment_affiliations(rows: list[dict]) -> int:
811782
for row in rows
812783
],
813784
)
814-
return len(rows)

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

Lines changed: 59 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,11 @@
1-
from pydantic import BaseModel
1+
from __future__ import annotations
2+
3+
import uuid
4+
from datetime import datetime
5+
from typing import Any
6+
7+
from loguru import logger
8+
from pydantic import BaseModel, TypeAdapter, ValidationError
29

310

411
class AffiliationContributor(BaseModel):
@@ -25,3 +32,54 @@ class AffiliationFile(BaseModel):
2532
class AffiliationParseOutput(BaseModel):
2633
affiliations: list[AffiliationInfoItem] | None = None
2734
error: str | None = None
35+
36+
37+
_SNAPSHOT_ADAPTER = TypeAdapter(list[AffiliationInfoItem])
38+
39+
40+
class RepoAffiliationRegistry(BaseModel):
41+
repo_id: str
42+
file_path: str | None = None
43+
file_hash: str | None = None
44+
status: str
45+
snapshot: list[AffiliationInfoItem] | None = None
46+
last_run_at: datetime | None = None
47+
48+
@classmethod
49+
def from_db(cls, db_data: dict[str, Any]) -> RepoAffiliationRegistry:
50+
row = db_data.copy()
51+
52+
for key, value in row.items():
53+
if value is not None and isinstance(value, uuid.UUID):
54+
row[key] = str(value)
55+
56+
field_mapping = {
57+
"repoId": "repo_id",
58+
"filePath": "file_path",
59+
"fileHash": "file_hash",
60+
"lastRunAt": "last_run_at",
61+
}
62+
for db_field, model_field in field_mapping.items():
63+
if db_field in row:
64+
row[model_field] = row.pop(db_field)
65+
66+
snapshot = row.get("snapshot")
67+
if snapshot is not None:
68+
row["snapshot"] = cls._parse_snapshot(snapshot)
69+
70+
return cls(**row)
71+
72+
@staticmethod
73+
def _parse_snapshot(snapshot) -> list[AffiliationInfoItem] | None:
74+
if isinstance(snapshot, dict) and "affiliations" in snapshot:
75+
snapshot = snapshot["affiliations"]
76+
try:
77+
return _SNAPSHOT_ADAPTER.validate_python(snapshot)
78+
except ValidationError as error:
79+
logger.warning(f"Invalid affiliation snapshot in registry, will re-parse: {error}")
80+
return None
81+
82+
def snapshot_for_db(self) -> list[dict] | None:
83+
if self.snapshot is None:
84+
return None
85+
return [item.model_dump() for item in self.snapshot]

0 commit comments

Comments
 (0)