From 121093877d84320a9c626dc25e08322031ad8928 Mon Sep 17 00:00:00 2001 From: Stas Moreinis Date: Fri, 21 Nov 2025 16:19:45 -0800 Subject: [PATCH 1/2] Update agents_acp_use_case.py --- .../domain/use_cases/agents_acp_use_case.py | 33 +++++++++++++++++-- 1 file changed, 30 insertions(+), 3 deletions(-) diff --git a/agentex/src/domain/use_cases/agents_acp_use_case.py b/agentex/src/domain/use_cases/agents_acp_use_case.py index 3f3c5a99..dd9f2194 100644 --- a/agentex/src/domain/use_cases/agents_acp_use_case.py +++ b/agentex/src/domain/use_cases/agents_acp_use_case.py @@ -1,9 +1,13 @@ +import asyncio import json from collections.abc import AsyncIterator, Callable from typing import Annotated, Any from fastapi import Depends +from src.adapters.authentication.exceptions import ( + AuthenticationServiceUnavailableError, +) from src.adapters.crud_store.exceptions import ItemDoesNotExist from src.api.schemas.authorization_types import ( AgentexResource, @@ -227,6 +231,31 @@ async def _execute_with_error_handling( await self.task_service.fail_task(task, str(e)) raise e + async def grant_with_retry(self, task: TaskEntity, max_retries: int = 3) -> None: + """Grant authorization for a task with retry""" + try: + await self.authorization_service.grant( + resource=AgentexResource.task(task.id), + ) + except AuthenticationServiceUnavailableError as e: + logger.error(f"Authentication service unavailable: {e}") + if max_retries > 0: + delay = 0.2 * (2 ** (max_retries - 1)) + logger.error( + f"Authentication service unavailable: {e}. Retrying in {delay}s..." + ) + await asyncio.sleep(delay) + return await self.grant_with_retry(task, max_retries - 1) + else: + logger.error( + f"Authentication service unavailable: {e}. Max retries reached." + ) + raise e from e + except Exception as e: + logger.error(f"Error granting authorization for task {task.id}: {e}") + await self.task_service.fail_task(task, str(e)) + raise e from e + async def _get_or_create_task( self, *, @@ -273,9 +302,7 @@ async def _get_or_create_task( agent=agent, task_name=task_name, task_params=task_params ) logger.info(f"[agent_id={agent.id}] Created task {task.id}") - await self.authorization_service.grant( - resource=AgentexResource.task(task.id), - ) + await self.grant_with_retry(task) return task async def handle_rpc_request( From febc6eb866a9fc969d9026b347196891da5a4e99 Mon Sep 17 00:00:00 2001 From: Stas Moreinis Date: Fri, 21 Nov 2025 17:50:18 -0800 Subject: [PATCH 2/2] Update agents_acp_use_case.py --- agentex/src/domain/use_cases/agents_acp_use_case.py | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/agentex/src/domain/use_cases/agents_acp_use_case.py b/agentex/src/domain/use_cases/agents_acp_use_case.py index dd9f2194..dee4d8cd 100644 --- a/agentex/src/domain/use_cases/agents_acp_use_case.py +++ b/agentex/src/domain/use_cases/agents_acp_use_case.py @@ -1,5 +1,6 @@ import asyncio import json +import random from collections.abc import AsyncIterator, Callable from typing import Annotated, Any @@ -231,21 +232,20 @@ async def _execute_with_error_handling( await self.task_service.fail_task(task, str(e)) raise e - async def grant_with_retry(self, task: TaskEntity, max_retries: int = 3) -> None: + async def grant_with_retry(self, task: TaskEntity, attempts: int = 0) -> None: """Grant authorization for a task with retry""" try: await self.authorization_service.grant( resource=AgentexResource.task(task.id), ) except AuthenticationServiceUnavailableError as e: - logger.error(f"Authentication service unavailable: {e}") - if max_retries > 0: - delay = 0.2 * (2 ** (max_retries - 1)) + if attempts < 3: + delay = 0.2 * (2**attempts) + random.uniform(0, 0.1) logger.error( f"Authentication service unavailable: {e}. Retrying in {delay}s..." ) await asyncio.sleep(delay) - return await self.grant_with_retry(task, max_retries - 1) + return await self.grant_with_retry(task, attempts + 1) else: logger.error( f"Authentication service unavailable: {e}. Max retries reached."