Coverage for chalicelib/utils/pg_client.py: 55%

165 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 12:56 +0000

1import logging 

2import time 

3from threading import Semaphore 

4 

5import psycopg2 

6import psycopg2.extras 

7from decouple import config 

8from psycopg2 import pool 

9 

10logger = logging.getLogger(__name__) 

11 

12_PG_CONFIG = {"host": config("pg_host"), 

13 "database": config("pg_dbname"), 

14 "user": config("pg_user"), 

15 "password": config("pg_password"), 

16 "port": config("pg_port", cast=int), 

17 "application_name": config("APP_NAME", default="PY")} 

18PG_CONFIG = dict(_PG_CONFIG) 

19if config("PG_TIMEOUT", cast=int, default=0) > 0: 19 ↛ 22line 19 didn't jump to line 22 because the condition on line 19 was always true

20 PG_CONFIG["options"] = f"-c statement_timeout={config('PG_TIMEOUT', cast=int) * 1000}" 

21 

22if config('PG_POOL', cast=bool, default=True): 22 ↛ 33line 22 didn't jump to line 33 because the condition on line 22 was always true

23 PG_CONFIG = { 

24 **PG_CONFIG, 

25 # Keepalive settings 

26 "keepalives": 1, # Enable keepalives 

27 "keepalives_idle": 300, # Seconds before sending keepalive 

28 "keepalives_interval": 10, # Seconds between keepalives 

29 "keepalives_count": 3 # Number of keepalives before giving up 

30 } 

31 

32 

33class ORThreadedConnectionPool(psycopg2.pool.ThreadedConnectionPool): 

34 def __init__(self, minconn, maxconn, *args, **kwargs): 

35 self._semaphore = Semaphore(maxconn) 

36 super().__init__(minconn, maxconn, *args, **kwargs) 

37 

38 def getconn(self, *args, **kwargs): 

39 self._semaphore.acquire() 

40 while True: 

41 try: 

42 conn = super().getconn(*args, **kwargs) 

43 if not config("PG_CHECK_CONNECTION", cast=bool, default=False): 43 ↛ 45line 43 didn't jump to line 45 because the condition on line 43 was always true

44 return conn 

45 try: 

46 with conn.cursor() as cur: 

47 cur.execute('SELECT 1') 

48 return conn 

49 except (psycopg2.InterfaceError, psycopg2.OperationalError): 

50 try: 

51 super().putconn(conn, close=True) 

52 except Exception: 

53 pass 

54 # Loop will retry 

55 except psycopg2.pool.PoolError as e: 

56 if str(e) == "connection pool is closed": 

57 make_pool() 

58 raise e 

59 

60 def putconn(self, *args, **kwargs): 

61 try: 

62 super().putconn(*args, **kwargs) 

63 self._semaphore.release() 

64 except psycopg2.pool.PoolError as e: 

65 if str(e) == "trying to put unkeyed connection": 

66 logger.warning("!!! trying to put unkeyed connection") 

67 logger.warning(f"env-PG_POOL:{config('PG_POOL', default=None)}") 

68 return 

69 raise e 

70 

71 

72postgreSQL_pool: ORThreadedConnectionPool = None 

73 

74RETRY_MAX = config("PG_RETRY_MAX", cast=int, default=50) 

75RETRY_INTERVAL = config("PG_RETRY_INTERVAL", cast=int, default=2) 

76RETRY = 0 

77 

78 

79def make_pool(): 

80 if not config('PG_POOL', cast=bool, default=True): 80 ↛ 81line 80 didn't jump to line 81 because the condition on line 80 was never true

81 logger.info("PG_POOL is disabled, not creating a new one") 

82 return 

83 global postgreSQL_pool 

84 global RETRY 

85 if postgreSQL_pool is not None: 85 ↛ 86line 85 didn't jump to line 86 because the condition on line 85 was never true

86 try: 

87 postgreSQL_pool.closeall() 

88 except (Exception, psycopg2.DatabaseError) as error: 

89 logger.error("Error while closing all connexions to PostgreSQL", exc_info=error) 

90 try: 

91 postgreSQL_pool = ORThreadedConnectionPool(config("PG_MINCONN", cast=int, default=4), 

92 config("PG_MAXCONN", cast=int, default=8), 

93 **PG_CONFIG) 

94 if postgreSQL_pool is not None: 94 ↛ exitline 94 didn't return from function 'make_pool' because the condition on line 94 was always true

95 logger.info("Connection pool created successfully") 

96 except (Exception, psycopg2.DatabaseError) as error: 

97 logger.error("Error while connecting to PostgreSQL", exc_info=error) 

98 if RETRY < RETRY_MAX: 

99 RETRY += 1 

100 logger.info(f"Waiting for {RETRY_INTERVAL}s before retry n°{RETRY}") 

101 time.sleep(RETRY_INTERVAL) 

102 make_pool() 

103 else: 

104 raise error 

105 

106 

107class PostgresClient: 

108 connection = None 

109 cursor = None 

110 long_query = False 

111 unlimited_query = False 

112 

113 def __init__(self, long_query=False, unlimited_query=False, use_pool=True): 

114 self.long_query = long_query 

115 self.unlimited_query = unlimited_query 

116 self.use_pool = use_pool 

117 if unlimited_query: 117 ↛ 118line 117 didn't jump to line 118 because the condition on line 117 was never true

118 long_config = dict(_PG_CONFIG) 

119 long_config["application_name"] += "-UNLIMITED" 

120 self.connection = psycopg2.connect(**long_config) 

121 elif long_query: 121 ↛ 122line 121 didn't jump to line 122 because the condition on line 121 was never true

122 long_config = dict(_PG_CONFIG) 

123 long_config["application_name"] += "-LONG" 

124 if config('PG_TIMEOUT_LONG', cast=int, default=1) > 0: 

125 long_config["options"] = f"-c statement_timeout=" \ 

126 f"{config('PG_TIMEOUT_LONG', cast=int, default=5 * 60) * 1000}" 

127 else: 

128 logger.info("Disabled timeout for long query") 

129 self.connection = psycopg2.connect(**long_config) 

130 elif not use_pool or not config('PG_POOL', cast=bool, default=True): 

131 single_config = dict(_PG_CONFIG) 

132 single_config["application_name"] += "-NOPOOL" 

133 if config('PG_TIMEOUT', cast=int, default=1) > 0: 133 ↛ 135line 133 didn't jump to line 135 because the condition on line 133 was always true

134 single_config["options"] = f"-c statement_timeout={config('PG_TIMEOUT', cast=int, default=30) * 1000}" 

135 self.connection = psycopg2.connect(**single_config) 

136 else: 

137 self.connection = postgreSQL_pool.getconn() 

138 

139 def __enter__(self): 

140 if self.cursor is None: 140 ↛ 145line 140 didn't jump to line 145 because the condition on line 140 was always true

141 self.cursor = self.connection.cursor(cursor_factory=psycopg2.extras.RealDictCursor) 

142 self.cursor.cursor_execute = self.cursor.execute 

143 self.cursor.execute = self.__execute 

144 self.cursor.recreate = self.recreate_cursor 

145 return self.cursor 

146 

147 def __exit__(self, *args): 

148 dead_connection = False 

149 try: 

150 self.connection.commit() 

151 self.cursor.close() 

152 if not self.use_pool or self.long_query or self.unlimited_query: 

153 self.connection.close() 

154 except Exception as error: 

155 logger.error("Error while committing/closing PG-connection", exc_info=error) 

156 dead_connection = isinstance(error, (psycopg2.OperationalError, psycopg2.InterfaceError)) \ 

157 or bool(self.connection.closed) 

158 if str(error) == "connection already closed" \ 

159 and self.use_pool \ 

160 and not self.long_query \ 

161 and not self.unlimited_query \ 

162 and config('PG_POOL', cast=bool, default=True): 

163 logger.info("Recreating the connexion pool") 

164 make_pool() 

165 else: 

166 raise error 

167 finally: 

168 if config('PG_POOL', cast=bool, default=True) \ 

169 and self.use_pool \ 

170 and not self.long_query \ 

171 and not self.unlimited_query: 

172 postgreSQL_pool.putconn(self.connection, close=dead_connection) 

173 elif dead_connection and not self.connection.closed: 173 ↛ 174line 173 didn't jump to line 174 because the condition on line 173 was never true

174 try: 

175 self.connection.close() 

176 except Exception as error: 

177 logger.error("Error while closing dead PG-connection", exc_info=error) 

178 

179 def __execute(self, query, vars=None): 

180 try: 

181 result = self.cursor.cursor_execute(query=query, vars=vars) 

182 except psycopg2.Error as error: 

183 logger.error(f"!!! Error of type:{type(error)} while executing query:") 

184 logger.error(query) 

185 logger.info("starting rollback to allow future execution") 

186 try: 

187 self.connection.rollback() 

188 except psycopg2.InterfaceError as e: 

189 logger.error("!!! Error while rollbacking connection", exc_info=e) 

190 logger.error("!!! Trying to recreate the cursor") 

191 self.recreate_cursor() 

192 raise error 

193 return result 

194 

195 def recreate_cursor(self, rollback=False): 

196 if rollback: 

197 try: 

198 self.connection.rollback() 

199 except Exception as error: 

200 logger.error("Error while rollbacking connection for recreation", exc_info=error) 

201 try: 

202 self.cursor.close() 

203 except Exception as error: 

204 logger.error("Error while closing cursor for recreation", exc_info=error) 

205 self.cursor = None 

206 return self.__enter__() 

207 

208 

209async def init(): 

210 logger.info(f">use PG_POOL:{config('PG_POOL', default=True)}") 

211 make_pool() 

212 

213 

214async def terminate(): 

215 global postgreSQL_pool 

216 if postgreSQL_pool is not None: 216 ↛ exitline 216 didn't return from function 'terminate' because the condition on line 216 was always true

217 try: 

218 postgreSQL_pool.closeall() 

219 logger.info("Closed all connexions to PostgreSQL") 

220 except (Exception, psycopg2.DatabaseError) as error: 

221 logger.error("Error while closing all connexions to PostgreSQL", exc_info=error)