@@ -30,50 +30,67 @@ class PoolError(Exception):
3030 pass
3131
3232class Pool :
33- plugins : Dict [str , Callable [[], Plugin ]]
34- instances : Dict [str , List [PoolPlugin ]]
33+ _plugins : Dict [str , Callable [[], Plugin ]]
34+ _instances : Dict [str , List [PoolPlugin ]]
35+ _locks : Dict [str , Lock ]
3536
3637 def __init__ (self , max_instances = 1 ):
3738 self .max_instances = max_instances
38- self .instances = {}
39- self .plugins = {}
40- self .count = {}
39+ self ._instances = {}
40+ self ._locks = {}
41+ self ._plugins = {}
4142
4243 def add (self , name , source : Callable [[], Plugin ]):
43- self .plugins [name ] = source
44- self .instances [name ] = []
45-
46- def find_available (self , name ):
47- entry = self .instances [name ]
44+ self ._plugins [name ] = source
45+ self ._instances [name ] = []
46+ self ._locks [name ] = Lock ()
47+
48+ def count (self , name ):
49+ self ._locks [name ].acquire ()
50+ n = len (self ._instances [name ])
51+ self ._locks [name ].release ()
52+ return n
53+
54+ def _find_available (self , name ):
55+ entry = self ._instances [name ]
4856 for instance in entry :
4957 if not instance .active :
5058 return instance .make_active ()
5159 return None
5260
5361 async def async_get (self , name , timeout = None ):
5462 start = time .time ()
55- entry = self .instances [name ]
56-
57- p = self .find_available (name )
58- if p is not None :
59- return p
60-
61- if len (entry ) < self .max_instances :
62- p = PoolPlugin (self .plugins [name ](), active = True )
63- entry .append (p )
64- self .instances [name ] = entry
65- return p
66-
67- while True :
68- p = self .find_available (name )
63+ entry = self ._instances [name ]
64+ lock = self ._locks [name ]
65+ lock .acquire ()
66+ try :
67+ p = self ._find_available (name )
6968 if p is not None :
69+ lock .release ()
7070 return p
71- else :
72- if timeout is None :
73- await asyncio .sleep (0 )
74- continue
75- elif (time .time () - start ) >= timeout :
76- raise PoolError ("Timed out getting instance for key " + name )
71+
72+ if len (entry ) < self .max_instances :
73+ p = PoolPlugin (self ._plugins [name ](), active = True )
74+ entry .append (p )
75+ self ._instances [name ] = entry
76+ lock .release ()
77+ return p
78+
79+ while True :
80+ p = self ._find_available (name )
81+ if p is not None :
82+ lock .release ()
83+ return p
84+ else :
85+ if timeout is None :
86+ await asyncio .sleep (0 )
87+ continue
88+ elif (time .time () - start ) >= timeout :
89+ lock .release ()
90+ raise PoolError ("Timed out getting instance for key " + name )
91+ except Exception as exc :
92+ lock .release ()
93+ raise exc
7794
7895 def get (self , name , timeout = None ):
7996 fut = self .async_get (name , timeout )
0 commit comments