#!/usr/bin/env python3 """ KeyCRM webhook demo stand (GuardLabs own test stand, not a client project). Flow: KeyCRM fires a webhook on new order / status change -> this service validates + dedupes the event -> sends a human-readable notification to Telegram (live public channel) -> maps the order to a SalesDrive lead payload (dry-run unless SALESDRIVE_URL is set) -> keeps a public event log for the demo page. Endpoints (behind nginx at /demo/keycrm-hook/api/): POST /api/webhook - accept KeyCRM-style webhook GET /api/log - last events (sanitized) GET /api/example - example payload for testing GET /api/health - liveness """ import json import os import re import threading import time import html from collections import deque from datetime import datetime, timezone, timedelta import requests from flask import Flask, request, jsonify BASE = os.path.dirname(os.path.abspath(__file__)) INFRA = json.load(open(os.path.join(BASE, "tg_infra.json"))) BOT_TOKEN = INFRA["bot_token"] CHANNEL = "@" + INFRA["channel_username"] CHANNEL_LINK = INFRA["channel_link"] SALESDRIVE_URL = os.environ.get("SALESDRIVE_URL", "") # empty -> dry-run LOG_FILE = os.path.join(BASE, "events.jsonl") KYIV = timezone(timedelta(hours=3)) MAX_BODY = 16 * 1024 # bytes RATE_PER_IP = 6 # webhooks per minute per IP TG_HOURLY_CAP = 40 # hard cap on telegram sends per hour (public demo) DEDUP_WINDOW = 600 # seconds app = Flask(__name__) lock = threading.Lock() events = deque(maxlen=50) # public log ring ip_hits = {} # ip -> [timestamps] seen_keys = {} # dedupe key -> timestamp tg_sends = deque(maxlen=200) # timestamps of tg sends started = time.time() counters = {"received": 0, "sent_tg": 0, "duplicates": 0, "rejected": 0} def now_iso(): return datetime.now(KYIV).strftime("%Y-%m-%d %H:%M:%S") def clean(s, limit=120): """Strip anything dangerous from user-supplied strings before they hit TG/log.""" if not isinstance(s, (str, int, float)): return "" s = str(s)[:limit] return re.sub(r"[\x00-\x1f<>]", " ", s).strip() def mask_phone(p): p = re.sub(r"[^\d+]", "", str(p))[:16] if len(p) >= 7: return p[:6] + "***" + p[-2:] return p def parse_order(payload): """ Accept both KeyCRM webhook envelope {"event": "...", "context": {...}} and a bare order object. Returns (event, order_dict) or (None, None). """ if not isinstance(payload, dict): return None, None if "context" in payload and isinstance(payload["context"], dict): return clean(payload.get("event", "order.created"), 60), payload["context"] if "grand_total" in payload or "buyer" in payload or "products" in payload: return "order.created", payload return None, None def order_fields(order): buyer = order.get("buyer") or order.get("client") or {} if not isinstance(buyer, dict): buyer = {} products = order.get("products") or [] if not isinstance(products, list): products = [] items = [] for p in products[:10]: if isinstance(p, dict): items.append({ "name": clean(p.get("name") or p.get("title") or "item", 80), "qty": int(p.get("quantity") or p.get("amount") or 1), "price": float(p.get("price") or p.get("price_sold") or 0), }) return { "id": clean(order.get("id") or order.get("source_uuid") or "?", 32), "name": clean(buyer.get("full_name") or buyer.get("name") or "невідомо", 80), "phone": mask_phone(buyer.get("phone") or ""), "email": clean(buyer.get("email") or "", 80), "total": float(order.get("grand_total") or order.get("total") or sum(i["qty"] * i["price"] for i in items)), "source": clean(order.get("source") or order.get("channel") or "API", 40), "status": clean(order.get("status") or order.get("status_name") or "new", 40), "manager_comment": clean(order.get("manager_comment") or "", 120), "items": items, } def to_salesdrive(f): """Map a KeyCRM order to a SalesDrive form-handler payload.""" name_parts = f["name"].split(" ", 1) return { "fName": name_parts[0], "lName": name_parts[1] if len(name_parts) > 1 else "", "phone": f["phone"], "email": f["email"], "comment": "KeyCRM order #%s (%s). %s" % (f["id"], f["source"], f["manager_comment"]), "products": [ {"name": i["name"], "costPerItem": i["price"], "amount": i["qty"]} for i in f["items"] ], "saleAmount": f["total"], "con_sourceComment": "keycrm-webhook-bridge", } def tg_message(event, f): lines = ["\U0001F6D2 Нове замовлення #%s" % f["id"]] if event and "status" in event: lines = ["\U0001F504 Замовлення #%s: зміна статусу" % f["id"]] lines.append("Клієнт: %s, %s" % (f["name"], f["phone"] or "тел. не вказано")) if f["items"]: for i in f["items"][:5]: lines.append(" · %s x%d = %.0f грн" % (i["name"], i["qty"], i["qty"] * i["price"])) lines.append("Сума: %.0f грн" % f["total"]) lines.append("Джерело: %s | Статус: %s" % (f["source"], f["status"])) lines.append("SalesDrive: %s" % ("відправлено" if SALESDRIVE_URL else "dry-run (мапінг готовий)")) return "\n".join(lines) def send_tg(text): now = time.time() while tg_sends and now - tg_sends[0] > 3600: tg_sends.popleft() if len(tg_sends) >= TG_HOURLY_CAP: return {"ok": False, "skipped": "hourly cap reached"} r = requests.post( "https://api.telegram.org/bot%s/sendMessage" % BOT_TOKEN, json={"chat_id": CHANNEL, "text": text}, timeout=10, ) d = r.json() if d.get("ok"): tg_sends.append(now) return {"ok": True, "message_id": d["result"]["message_id"], "channel": CHANNEL} return {"ok": False, "error": clean(str(d.get("description")), 100)} def push_event(entry): with lock: events.appendleft(entry) try: with open(LOG_FILE, "a") as fh: fh.write(json.dumps(entry, ensure_ascii=False) + "\n") except OSError: pass def rate_ok(ip): now = time.time() hits = [t for t in ip_hits.get(ip, []) if now - t < 60] if len(hits) >= RATE_PER_IP: ip_hits[ip] = hits return False hits.append(now) ip_hits[ip] = hits return True @app.route("/api/webhook", methods=["POST"]) def webhook(): counters["received"] += 1 ip = request.headers.get("X-Real-IP", request.remote_addr or "?") if request.content_length and request.content_length > MAX_BODY: counters["rejected"] += 1 return jsonify({"ok": False, "error": "payload too large"}), 413 if not rate_ok(ip): counters["rejected"] += 1 return jsonify({"ok": False, "error": "rate limit: 6 req/min per IP"}), 429 try: payload = request.get_json(force=True, silent=False) except Exception: counters["rejected"] += 1 return jsonify({"ok": False, "error": "invalid JSON"}), 400 event, order = parse_order(payload) if order is None: counters["rejected"] += 1 return jsonify({"ok": False, "error": "unrecognized payload: expected KeyCRM envelope {event, context} or bare order"}), 422 f = order_fields(order) # Dedupe: KeyCRM retries webhooks on non-200; same event+order within window is a duplicate. key = "%s|%s|%s" % (event, f["id"], f["status"]) now = time.time() for k in [k for k, t in seen_keys.items() if now - t > DEDUP_WINDOW]: del seen_keys[k] if key in seen_keys: counters["duplicates"] += 1 push_event({"ts": now_iso(), "event": event, "order_id": f["id"], "result": "duplicate (skipped resend)"}) return jsonify({"ok": True, "duplicate": True}), 200 seen_keys[key] = now tg_result = send_tg(tg_message(event, f)) if tg_result.get("ok"): counters["sent_tg"] += 1 sd_payload = to_salesdrive(f) sd_result = {"mode": "dry-run", "mapped": True} if SALESDRIVE_URL: try: sd = requests.post(SALESDRIVE_URL, json=sd_payload, timeout=10) sd_result = {"mode": "live", "status": sd.status_code} except requests.RequestException as e: sd_result = {"mode": "live", "error": clean(str(e), 100)} entry = { "ts": now_iso(), "event": event, "order_id": f["id"], "client": f["name"], "phone": f["phone"], "total": f["total"], "source": f["source"], "status": f["status"], "telegram": tg_result, "salesdrive": {"payload": sd_payload, "result": sd_result}, "from_ip": ip.rsplit(".", 1)[0] + ".x" if "." in ip else "?", } push_event(entry) return jsonify({"ok": True, "telegram": tg_result, "salesdrive": sd_result, "channel": CHANNEL_LINK}), 200 @app.route("/api/log") def log_view(): with lock: return jsonify({"events": list(events), "counters": counters}) @app.route("/api/example") def example(): return jsonify({ "event": "order.created", "context": { "id": 1024, "grand_total": 2450, "source": "Instagram Shop", "status": "new", "buyer": {"full_name": "Олена Петренко", "phone": "+380671234567", "email": "olena@example.com"}, "products": [ {"name": "Керамічна ваза Terra", "quantity": 1, "price": 1650}, {"name": "Набір свічок Soy Mini", "quantity": 2, "price": 400} ] } }) @app.route("/api/health") def health(): return jsonify({"ok": True, "uptime_s": int(time.time() - started), "counters": counters, "channel": CHANNEL_LINK}) if __name__ == "__main__": app.run(host="127.0.0.1", port=8201, threaded=True)