Skip to content

Commit fd0a7c7

Browse files
committed
refactor: retry malformed affiliation parses once before unusable
Signed-off-by: Yeganathan S <63534555+skwowet@users.noreply.github.com>
1 parent fd55065 commit fd0a7c7

1 file changed

Lines changed: 21 additions & 9 deletions

File tree

services/apps/git_integration/src/crowdgit/services/affiliation/affiliation_service.py

Lines changed: 21 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77

88
import aiofiles
99
import aiofiles.os
10+
from pydantic import ValidationError
1011

1112
from crowdgit.database.crud import (
1213
fetch_member_organizations,
@@ -459,11 +460,25 @@ def normalize_parsed_affiliations(
459460

460461
async def parse_affiliations(self, content: str) -> tuple[list[AffiliationInfoItem], float]:
461462
"""Extract affiliations with AI, splitting large files into chunks when needed."""
463+
464+
async def invoke_parse(file_content: str):
465+
for attempt in range(2):
466+
try:
467+
return await invoke_bedrock(
468+
self.get_extraction_prompt(file_content),
469+
pydantic_model=AffiliationParseOutput,
470+
)
471+
except ValidationError:
472+
if attempt == 0:
473+
self.logger.info("Malformed affiliation parse response, retrying once")
474+
continue
475+
raise AffiliationAnalysisError(
476+
retain_file_hash=True,
477+
error_message="Affiliation file could not be parsed cleanly after retry",
478+
) from None
479+
462480
if len(content) <= self.MAX_CHUNK_SIZE:
463-
parse_result = await invoke_bedrock(
464-
self.get_extraction_prompt(content),
465-
pydantic_model=AffiliationParseOutput,
466-
)
481+
parse_result = await invoke_parse(content)
467482

468483
affiliations = parse_result.output.affiliations
469484
if affiliations is not None:
@@ -504,10 +519,7 @@ async def parse_affiliations(self, content: str) -> tuple[list[AffiliationInfoIt
504519

505520
async def process_chunk(chunk: str):
506521
async with semaphore:
507-
return await invoke_bedrock(
508-
self.get_extraction_prompt(chunk),
509-
pydantic_model=AffiliationParseOutput,
510-
)
522+
return await invoke_parse(chunk)
511523

512524
chunk_results = await asyncio.gather(*[process_chunk(chunk) for chunk in chunks])
513525

@@ -807,7 +819,7 @@ async def process_affiliations(
807819
)
808820
)
809821

810-
self.logger.info("Starting affiliations")
822+
self.logger.info(f"Starting affiliations processing for repo: {batch_info.remote}")
811823

812824
saved_file_path = registry.file_path if registry else None
813825
latest_file_path, discovery_cost = await self.resolve_affiliation_file(

0 commit comments

Comments
 (0)