66from collections import defaultdict
77from datetime import datetime , timedelta , timezone
88
9+ import aiohttp
910import escapism
10- import requests
1111from yarl import URL
1212
1313from .cache import ttl_lru_cache
2323prometheus_password = os .environ .get ("PROMETHEUS_PASSWORD" , "" )
2424
2525
26- def query_prometheus (query : str , date_range : DateRange , step : str ) -> requests .Response :
26+ async def query_prometheus (
27+ client : aiohttp .ClientSession , query : str , date_range : DateRange , step : str
28+ ) -> dict :
2729 """
2830 Query the Prometheus server with the given query over a date range.
2931
@@ -42,9 +44,7 @@ def query_prometheus(query: str, date_range: DateRange, step: str) -> requests.R
4244 scheme = "http" , host = prometheus_host , port = prometheus_port
4345 )
4446 if prometheus_username != "" and prometheus_password != "" :
45- prometheus_auth = requests .auth .HTTPBasicAuth (
46- prometheus_username , prometheus_password
47- )
47+ prometheus_auth = aiohttp .BasicAuth (prometheus_username , prometheus_password )
4848 else :
4949 prometheus_auth = None
5050 parameters = {
@@ -53,15 +53,16 @@ def query_prometheus(query: str, date_range: DateRange, step: str) -> requests.R
5353 "end" : to_date ,
5454 "step" : step ,
5555 }
56- query_api = URL ( prometheus_api .with_path ("/api/v1/query_range" ))
57- with requests .get (query_api , params = parameters , auth = prometheus_auth ) as response :
56+ query_api = prometheus_api .with_path ("/api/v1/query_range" ). with_query ( parameters )
57+ async with client .get (query_api , auth = prometheus_auth ) as response :
5858 logger .info (f"Querying Prometheus: { response .url } " )
5959 response .raise_for_status ()
60- result = response .json ()
60+ result = await response .json ()
6161 return result
6262
6363
64- def query_usage (
64+ async def query_usage (
65+ client : aiohttp .ClientSession ,
6566 date_range : DateRange ,
6667 hub_name : str | None ,
6768 component_name : str | None ,
@@ -84,23 +85,18 @@ def query_usage(
8485 if component_name is None :
8586 # Query all components defined in USAGE_MAP
8687 for component , params in USAGE_MAP .items ():
87- try :
88- response = query_prometheus (
89- params ["query" ], date_range , step = params ["step" ]
90- )
91- except requests .exceptions .RequestException :
92- raise
88+ response = await query_prometheus (
89+ client , params ["query" ], date_range , step = params ["step" ]
90+ )
9391 result .extend (_process_response (response , component ))
9492 else :
9593 # Query specific component only
96- try :
97- response = query_prometheus (
98- USAGE_MAP [component_name ]["query" ],
99- date_range ,
100- step = USAGE_MAP [component_name ]["step" ],
101- )
102- except requests .exceptions .RequestException :
103- raise
94+ response = await query_prometheus (
95+ client ,
96+ USAGE_MAP [component_name ]["query" ],
97+ date_range ,
98+ step = USAGE_MAP [component_name ]["step" ],
99+ )
104100 result .extend (_process_response (response , component_name ))
105101 # Calculate daily cost factors from absolute usage totals)
106102 result = _calculate_daily_cost_factors (result , hub_name = hub_name )
@@ -111,9 +107,9 @@ def query_usage(
111107
112108
113109def _process_response (
114- response : requests . Response ,
110+ response : dict ,
115111 component_name : str ,
116- ) -> dict :
112+ ) -> list [ dict ] :
117113 """
118114 Process the response from the Prometheus server to extract absolute usage data.
119115
@@ -256,7 +252,8 @@ def _calculate_daily_cost_factors(
256252
257253
258254@ttl_lru_cache (seconds_to_live = 3600 )
259- def query_user_groups (
255+ async def query_user_groups (
256+ client : aiohttp .ClientSession ,
260257 hub_name : str | None = None ,
261258 user_name : str | None = None ,
262259 group_name : str | None = None ,
@@ -266,17 +263,13 @@ def query_user_groups(
266263 """
267264 now_date = get_now_date () - timedelta (days = 1 )
268265 date_range = DateRange (start_date = now_date , end_date = now_date )
269- try :
270- response = query_prometheus (USER_GROUP_INFO , date_range , step = "1d" )
271- except requests .exceptions .RequestException as e :
272- logger .exception (f"HTTP request failed: { e } " )
273- raise
266+ response = await query_prometheus (client , USER_GROUP_INFO , date_range , step = "1d" )
274267 result = _process_user_groups (response , hub_name , user_name , group_name )
275268 return result
276269
277270
278271def _process_user_groups (
279- response : requests . Response ,
272+ response : dict ,
280273 hub_name : str | None = None ,
281274 user_name : str | None = None ,
282275 group_name : str | None = None ,
@@ -306,16 +299,13 @@ def _process_user_groups(
306299
307300
308301@ttl_lru_cache (seconds_to_live = 3600 )
309- def query_users_with_multiple_groups (
302+ async def query_users_with_multiple_groups (
303+ client : aiohttp .ClientSession ,
310304 date_range : DateRange ,
311305 hub_name : str | None = None ,
312306 user_name : str | None = None ,
313307) -> list [dict ]:
314- try :
315- response = query_user_groups (hub_name = hub_name , user_name = user_name )
316- except requests .exceptions .RequestException as e :
317- logger .exception (f"HTTP request failed: { e } " )
318- raise
308+ response = await query_user_groups (client , hub_name = hub_name , user_name = user_name )
319309 grouped = defaultdict (
320310 lambda : {"username" : None , "hub" : None , "usergroups" : [], "has_multiple" : False }
321311 )
@@ -340,16 +330,13 @@ def query_users_with_multiple_groups(
340330
341331
342332@ttl_lru_cache (seconds_to_live = 3600 )
343- def query_users_with_no_groups (
333+ async def query_users_with_no_groups (
334+ client : aiohttp .ClientSession ,
344335 date_range : DateRange ,
345336 hub_name : str | None = None ,
346337 user_name : str | None = None ,
347338) -> list [dict ]:
348- try :
349- response = query_user_groups (hub_name = hub_name , user_name = user_name )
350- except requests .exceptions .RequestException as e :
351- logger .exception (f"HTTP request failed: { e } " )
352- raise
339+ response = await query_user_groups (client , hub_name = hub_name , user_name = user_name )
353340 grouped = defaultdict (lambda : {"username" : None , "hub" : None })
354341 for entry in response :
355342 key = (entry ["username" ], entry ["hub" ])
0 commit comments