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
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:56 +0000
1import logging
2from typing import Optional
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
10logger = logging.getLogger(__name__)
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
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
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
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
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
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
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"]
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)
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"}}
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']}")
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"]}
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