Coverage for chalicelib/core/webhook.py: 31%

109 statements  

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

1import logging 

2from typing import Optional 

3 

4import schemas 

5from chalicelib.utils import pg_client, helper, ssrf 

6from chalicelib.utils.TimeUTC import TimeUTC 

7from chalicelib.utils.log import sanitize 

8from fastapi import HTTPException, status 

9 

10logger = logging.getLogger(__name__) 

11 

12 

13def get_by_id(webhook_id): 

14 with pg_client.PostgresClient() as cur: 

15 cur.execute( 

16 cur.mogrify("""\ 

17 SELECT w.* 

18 FROM public.webhooks AS w  

19 WHERE w.webhook_id =%(webhook_id)s AND deleted_at ISNULL;""", 

20 {"webhook_id": webhook_id}) 

21 ) 

22 w = helper.dict_to_camel_case(cur.fetchone()) 

23 if w: 

24 w["createdAt"] = TimeUTC.datetime_to_timestamp(w["createdAt"]) 

25 return w 

26 

27 

28def get_webhook(tenant_id, webhook_id, webhook_type='webhook'): 

29 with pg_client.PostgresClient() as cur: 

30 cur.execute( 

31 cur.mogrify("""SELECT w.* 

32 FROM public.webhooks AS w  

33 WHERE w.webhook_id =%(webhook_id)s  

34 AND deleted_at ISNULL AND type=%(webhook_type)s;""", 

35 {"webhook_id": webhook_id, "webhook_type": webhook_type}) 

36 ) 

37 w = helper.dict_to_camel_case(cur.fetchone()) 

38 if w: 38 ↛ 39line 38 didn't jump to line 39 because the condition on line 38 was never true

39 w["createdAt"] = TimeUTC.datetime_to_timestamp(w["createdAt"]) 

40 return w 

41 

42 

43def get_by_type(tenant_id, webhook_type): 

44 with pg_client.PostgresClient() as cur: 

45 cur.execute( 

46 cur.mogrify("""SELECT w.webhook_id,w.endpoint,w.auth_header,w.type,w.index,w.name,w.created_at 

47 FROM public.webhooks AS w  

48 WHERE w.type =%(type)s AND deleted_at ISNULL;""", 

49 {"type": webhook_type}) 

50 ) 

51 webhooks = helper.list_to_camel_case(cur.fetchall()) 

52 for w in webhooks: 52 ↛ 53line 52 didn't jump to line 53 because the loop on line 52 never started

53 w["createdAt"] = TimeUTC.datetime_to_timestamp(w["createdAt"]) 

54 return webhooks 

55 

56 

57def get_by_tenant(tenant_id, replace_none=False): 

58 with pg_client.PostgresClient() as cur: 

59 cur.execute("""SELECT w.* 

60 FROM public.webhooks AS w  

61 WHERE deleted_at ISNULL;""") 

62 all = helper.list_to_camel_case(cur.fetchall()) 

63 for w in all: 63 ↛ 64line 63 didn't jump to line 64 because the loop on line 63 never started

64 w["createdAt"] = TimeUTC.datetime_to_timestamp(w["createdAt"]) 

65 return all 

66 

67 

68def update(tenant_id, webhook_id, changes, replace_none=False): 

69 allow_update = ["name", "index", "authHeader", "endpoint"] 

70 with pg_client.PostgresClient() as cur: 

71 sub_query = [f"{helper.key_to_snake_case(k)} = %({k})s" for k in changes.keys() if k in allow_update] 

72 cur.execute( 

73 cur.mogrify(f"""\ 

74 UPDATE public.webhooks 

75 SET {','.join(sub_query)} 

76 WHERE webhook_id =%(id)s AND deleted_at ISNULL 

77 RETURNING *;""", 

78 {"id": webhook_id, **changes}) 

79 ) 

80 w = helper.dict_to_camel_case(cur.fetchone()) 

81 if w is None: 

82 raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail=f"webhook not found.") 

83 w["createdAt"] = TimeUTC.datetime_to_timestamp(w["createdAt"]) 

84 if replace_none: 

85 for k in w.keys(): 

86 if w[k] is None: 

87 w[k] = '' 

88 return w 

89 

90 

91def add(tenant_id, endpoint, auth_header=None, webhook_type='webhook', name="", replace_none=False): 

92 with pg_client.PostgresClient() as cur: 

93 query = cur.mogrify("""\ 

94 INSERT INTO public.webhooks(endpoint,auth_header,type,name) 

95 VALUES (%(endpoint)s, %(auth_header)s, %(type)s,%(name)s) 

96 RETURNING *;""", 

97 {"endpoint": endpoint, "auth_header": auth_header, 

98 "type": webhook_type, "name": name}) 

99 cur.execute( 

100 query 

101 ) 

102 w = helper.dict_to_camel_case(cur.fetchone()) 

103 w["createdAt"] = TimeUTC.datetime_to_timestamp(w["createdAt"]) 

104 if replace_none: 

105 for k in w.keys(): 

106 if w[k] is None: 

107 w[k] = '' 

108 return w 

109 

110 

111def exists_by_name(name: str, exclude_id: Optional[int], webhook_type: str = schemas.WebhookType.WEBHOOK, 

112 tenant_id: Optional[int] = None) -> bool: 

113 with pg_client.PostgresClient() as cur: 

114 query = cur.mogrify(f"""SELECT EXISTS(SELECT 1  

115 FROM public.webhooks 

116 WHERE name ILIKE %(name)s 

117 AND deleted_at ISNULL 

118 AND type=%(webhook_type)s 

119 {"AND webhook_id!=%(exclude_id)s" if exclude_id else ""}) AS exists;""", 

120 {"name": name, "exclude_id": exclude_id, "webhook_type": webhook_type}) 

121 cur.execute(query) 

122 row = cur.fetchone() 

123 return row["exists"] 

124 

125 

126def add_edit(tenant_id, data: schemas.WebhookSchema, replace_none=None): 

127 if len(data.name) > 0 \ 

128 and exists_by_name(name=data.name, exclude_id=data.webhook_id): 

129 raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=f"name already exists.") 

130 if data.webhook_id is not None: 

131 return update(tenant_id=tenant_id, webhook_id=data.webhook_id, 

132 changes={"endpoint": data.endpoint, 

133 "authHeader": data.auth_header, 

134 "name": data.name}, 

135 replace_none=replace_none) 

136 else: 

137 return add(tenant_id=tenant_id, 

138 endpoint=data.endpoint, 

139 auth_header=data.auth_header, 

140 name=data.name, 

141 replace_none=replace_none) 

142 

143 

144def delete(tenant_id, webhook_id): 

145 with pg_client.PostgresClient() as cur: 

146 cur.execute( 

147 cur.mogrify("""\ 

148 UPDATE public.webhooks 

149 SET deleted_at = (now() at time zone 'utc') 

150 WHERE webhook_id =%(id)s AND deleted_at ISNULL 

151 RETURNING *;""", 

152 {"id": webhook_id}) 

153 ) 

154 return {"data": {"state": "success"}} 

155 

156 

157def trigger_batch(data_list): 

158 webhooks_map = {} 

159 for w in data_list: 

160 if w["destination"] not in webhooks_map: 

161 webhooks_map[w["destination"]] = get_by_id(webhook_id=w["destination"]) 

162 if webhooks_map[w["destination"]] is None: 

163 logger.error(f"!!Error webhook not found: webhook_id={w['destination']}") 

164 else: 

165 try: 

166 __trigger(hook=webhooks_map[w["destination"]], data=w["data"]) 

167 except Exception: 

168 logger.exception(f"!!Error while triggering webhook_id={w['destination']}") 

169 

170 

171def __trigger(hook, data): 

172 if hook is not None and hook["type"] == 'webhook': 

173 headers = {} 

174 if hook["authHeader"] is not None and len(hook["authHeader"]) > 0: 

175 headers = {"Authorization": hook["authHeader"]} 

176 

177 r = ssrf.post_json(endpoint=hook["endpoint"], json_data=data, headers=headers) 

178 if r.status_code != 200: 

179 logger.error("=======> webhook: something went wrong for:") 

180 logger.error(hook) 

181 logger.error(r.status_code) 

182 logger.error(sanitize(r.text)) 

183 return 

184 response = None 

185 try: 

186 response = r.json() 

187 except: 

188 try: 

189 response = r.text 

190 except: 

191 logger.info("no response found") 

192 return response