Skip to content

Commit ca661de

Browse files
committed
refactor OBPClient
1 parent 535d16c commit ca661de

7 files changed

Lines changed: 236 additions & 220 deletions

File tree

.github/instructions/ai.md.instructions.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,6 @@ applyTo: '**'
1919
- Keep documentation as simple as possible
2020

2121
# Personality
22-
- Avoid ecessive sycophancy and pandering. Be professional and respectful.
22+
- Avoid ecessive sycophancy and pandering.
2323
- Please disagree explicitly if you think the user is wrong or there is a better way to do something.
2424
- Do not always agree with the user.

src/auth/admin_client.py

Lines changed: 29 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -71,8 +71,29 @@ async def initialize(
7171

7272
except Exception as e:
7373
logger.error(f'❌ Failed to initialize admin OBP client: {e}')
74+
# Clean up any partially created resources
75+
await self._cleanup_partial_init()
7476
raise ValueError(f'Admin client initialization failed: {e}') from e
7577

78+
async def _cleanup_partial_init(self) -> None:
79+
"""Clean up resources from a failed initialization."""
80+
if self._client:
81+
try:
82+
await self._client.close()
83+
except Exception as cleanup_error:
84+
logger.debug(f'Error during cleanup: {cleanup_error}')
85+
86+
if self._auth and hasattr(self._auth, 'async_requests_client'):
87+
if self._auth.async_requests_client:
88+
try:
89+
await self._auth.async_requests_client.close()
90+
except Exception as cleanup_error:
91+
logger.debug(f'Error during cleanup: {cleanup_error}')
92+
93+
self._client = None
94+
self._auth = None
95+
self._initialized = False
96+
7697
def get_client(self) -> OBPClient:
7798
"""
7899
Get the singleton admin OBP client instance.
@@ -112,10 +133,16 @@ def is_initialized(self) -> bool:
112133

113134
async def close(self) -> None:
114135
"""Clean up resources during app shutdown."""
136+
# Close OBP client session
137+
if self._client:
138+
await self._client.close()
139+
logger.info('🔌 Admin OBP client session closed')
140+
141+
# Close auth session
115142
if self._auth and hasattr(self._auth, 'async_requests_client'):
116143
if self._auth.async_requests_client:
117144
await self._auth.async_requests_client.close()
118-
logger.info('🔌 Admin client HTTP session closed')
145+
logger.info('🔌 Admin auth HTTP session closed')
119146

120147
self._initialized = False
121148
self._client = None
@@ -160,7 +187,7 @@ def get_admin_client() -> OBPClient:
160187
>>>
161188
>>> # Later, anywhere in the app
162189
>>> admin_client = get_admin_client()
163-
>>> response = await admin_client.get("/obp/v6.0.0/banks")
190+
>>> banks = await admin_client.get("/obp/v6.0.0/banks")
164191
"""
165192
return _admin_manager.get_client()
166193

src/auth/auth.py

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -459,17 +459,15 @@ async def create_admin_direct_login_auth(
459459
version = os.getenv('OBP_API_VERSION', 'v6.0.0')
460460

461461
try:
462-
response_json = await obp_client.get(
462+
entitlements_response = await obp_client.get(
463463
f"/obp/{version}/my/entitlements"
464464
)
465+
entitlements_data = entitlements_response.json()
465466

466-
if not response_json:
467+
if not entitlements_data:
467468
logger.warning('Failed to fetch entitlements - received empty response')
468469
return admin_auth
469470

470-
# Parse the JSON string response
471-
entitlements_response = json.loads(response_json) if isinstance(response_json, str) else response_json
472-
473471
# Use provided entitlements or default set
474472
if required_entitlements is None:
475473
required_entitlements = [
@@ -479,7 +477,7 @@ async def create_admin_direct_login_auth(
479477
'CanGetSystemLevelDynamicEntities'
480478
]
481479

482-
all_present, missing = _check_entitlements(entitlements_response, required_entitlements)
480+
all_present, missing = _check_entitlements(entitlements_data, required_entitlements)
483481

484482
if not all_present:
485483
logger.warning(
@@ -490,5 +488,8 @@ async def create_admin_direct_login_auth(
490488
except Exception as e:
491489
logger.error(f'Failed to verify admin entitlements: {e}')
492490
# Don't raise - auth still works, just couldn't verify entitlements
491+
finally:
492+
# Always close the temporary client session
493+
await obp_client.close()
493494

494495
return admin_auth

src/checkpointer/entities.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ async def read(self, entity_id: str) -> dict:
8888
response = await self.client.get(
8989
path=f"{self.endpoint_url}/{entity_id}"
9090
)
91-
return json.loads(response) if isinstance(response, str) else response
91+
return response.json()
9292

9393
opey_checkpoint_entity = {
9494
"hasPersonalEntity": True,

src/checkpointer/obp_checkpoint_saver.py

Lines changed: 37 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,18 +10,50 @@
1010
)
1111
from langgraph.checkpoint.serde.jsonplus import JsonPlusSerializer
1212

13-
from src.client import obp_client
13+
from src.client.obp_client import OBPClient
1414
from src.checkpointer.entities import OpeyCheckpointEntity, OpeyCheckpointWriteEntity
15-
15+
from src.auth.admin_client import get_admin_client
1616
class OBPCheckpointSaver(BaseCheckpointSaver):
1717
"""Checkpoint saver using OBP Dynamic Entitites as a storage backend."""
1818

1919
def __init__(self, client: OBPClient):
2020
super().__init__()
21+
if not client:
22+
try:
23+
client = get_admin_client()
24+
except RuntimeError as e:
25+
raise ValueError("An OBPClient must be provided or the admin client must be initialized.") from e
26+
2127
self.client = client
2228
self.is_setup = False
2329
self.serde = JsonPlusSerializer()
2430

31+
async def _check_existing_setup(self) -> bool:
32+
"""
33+
Check if the required Dynamic Entities are already set up in OBP.
34+
35+
Returns:
36+
bool: True if setup exists, False otherwise.
37+
"""
38+
try:
39+
response = await self.client.get("/obp/v6.0.0/management/system-dynamic-entities")
40+
response_data = response.json()
41+
except Exception as e:
42+
raise RuntimeError(f"Error checking existing OBP setup: {e}") from e
43+
44+
if not response_data or not isinstance(response_data, dict):
45+
return False
46+
47+
existing_entities = response_data.get("dynamic_entities", [])
48+
required_entities = {OpeyCheckpointEntity.obp_entity_name(), OpeyCheckpointWriteEntity.obp_entity_name()}
49+
50+
for entity in existing_entities:
51+
entity_name = next(iter(entity))
52+
if entity_name in required_entities:
53+
required_entities.remove(entity_name)
54+
55+
return len(required_entities) == 0
56+
2557
def setup(self) -> None:
2658
"""
2759
Setup the Dynamic Entities in OBP if they do not exist.
@@ -30,6 +62,9 @@ def setup(self) -> None:
3062
if self.is_setup:
3163
return
3264

65+
66+
67+
3368
def put(
3469
self,
3570
config: RunnableConfig,

0 commit comments

Comments
 (0)