commit d19ec3e3a3e1cd384672f71907253a634776b6b8 Author: yjca Date: Tue Jul 7 14:28:24 2026 +0800 first commit diff --git a/.env b/.env new file mode 100644 index 0000000..220649f --- /dev/null +++ b/.env @@ -0,0 +1,11 @@ +# Aliyun SMS config +ALIBABA_CLOUD_ACCESS_KEY_ID=LTAI5tCpEwwAbxbqEM5rQDA9 +ALIBABA_CLOUD_ACCESS_KEY_SECRET=BUcLxYYmvsZhBnbqrq9tHjICWzDzrc +ALIBABA_CLOUD_REGION_ID=cn-hangzhou +ALIYUN_SMS_SIGN_NAME=浙江云览数字 +ALIYUN_SMS_TEMPLATE_CODE=SMS_508915141 +SMS_DB_PATH=data/sms_receivers.db +SMS_DB_TABLE=sms_receivers + +# true: preview only, false: actually send SMS +LOG_ONLY_NO_SMS=false diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..2ae2839 --- /dev/null +++ b/.gitignore @@ -0,0 +1 @@ +pass diff --git a/README.md b/README.md new file mode 100644 index 0000000..f562336 --- /dev/null +++ b/README.md @@ -0,0 +1,175 @@ +# 回调告警转短信服务 + +本项目是一个 Python 后端服务,用于接收 AI 视频分析平台的 HTTP 回调,并将关键告警信息转为阿里云短信发送给手机号列表。 + +## 启动指令 +pip install -r requirements.txt +py -m uvicorn app.main:app --host 0.0.0.0 --port 2336 + +核心目标: +1. 接收并解析上游回调 JSON。 +2. 生成短信模板变量 ACTION。 +3. 查询 SQLite 中已启用手机号并群发。 +4. 记录可追踪日志,返回标准化接口结果。 + +## 1. 功能概览 + +服务提供两个接口: +1. GET /healthz:健康检查。 +2. POST /callback:接收告警回调并触发短信逻辑。 + +POST /callback 的处理流程: +1. 解析 JSON(支持常见编码容错)。 +2. 提取事件标识 event_id(优先 snowflake_id,其次 analysis_job_id)。 +3. 脱敏写日志(图片 base64、密钥等敏感字段不落盘)。 +4. 生成短信预览文本:发现有${ACTION}行为,请去摄像头查看。 +5. 查询 SQLite 中 is_enabled=1 的手机号。 +6. 根据 LOG_ONLY_NO_SMS 决定仅预览还是实际调用阿里云接口。 +7. 返回统一结构:code、message、data。 + +## 2. 环境准备 + +在项目根目录执行: + +```bash +pip install -r requirements.txt +``` + +推荐 Python 3.11+。 + +## 3. 配置说明 + +项目默认通过 .env 读取配置,主要参数如下: + +1. ALIBABA_CLOUD_ACCESS_KEY_ID:阿里云 AK。 +2. ALIBABA_CLOUD_ACCESS_KEY_SECRET:阿里云 SK。 +3. ALIBABA_CLOUD_REGION_ID:区域,默认 cn-hangzhou。 +4. ALIYUN_SMS_SIGN_NAME:已审核通过的短信签名。 +5. ALIYUN_SMS_TEMPLATE_CODE:短信模板编码。 +6. SMS_DB_PATH:SQLite 数据库路径,默认 data/sms_receivers.db。 +7. SMS_DB_TABLE:手机号表名,默认 sms_receivers。 +8. LOG_ONLY_NO_SMS:是否仅预览不实发。 + +LOG_ONLY_NO_SMS 取值建议: +1. true:联调阶段使用,不调用阿里云发送。 +2. false:生产或实发测试使用,会真实发送短信。 + +## 4. 手机号来源与优先级 + +服务启动时会自动创建手机号表(如果不存在): + +```sql +CREATE TABLE IF NOT EXISTS sms_receivers ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + phone_number TEXT NOT NULL UNIQUE, + is_enabled INTEGER NOT NULL DEFAULT 1, + updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP +); +``` + +接收方选择规则: +1. 仅使用 SQLite 中 is_enabled=1 的手机号。 + +这意味着你在服务运行期间修改数据库,下一次 callback 会立即生效,无需重启。 + +常用维护 SQL: + +```sql +INSERT INTO sms_receivers (phone_number, is_enabled) VALUES ('17394641215', 1); +UPDATE sms_receivers SET is_enabled = 0, updated_at = CURRENT_TIMESTAMP WHERE phone_number = '17394641215'; +DELETE FROM sms_receivers WHERE phone_number = '17394641215'; +``` + +## 5. 启动服务 + +在项目根目录执行: + +```bash +py -m uvicorn app.main:app --host 0.0.0.0 --port 2336 +``` + +启动成功后会看到: +1. sms_db_ready +2. Uvicorn running on http://0.0.0.0:2336 + +## 6. 如何测试 + +### 6.1 健康检查 + +```bash +curl http://127.0.0.1:2336/healthz +``` + +### 6.2 回调测试 + +方式一:使用示例文件 test_callback_payload.json。 + +```bash +curl -X POST "http://127.0.0.1:2336/callback" \ + -H "Content-Type: application/json; charset=utf-8" \ + --data-binary "@test_callback_payload.json" +``` + +方式二:直接发最小 JSON。 + +```bash +curl -X POST "http://127.0.0.1:2336/callback" \ + -H "Content-Type: application/json" \ + -d "{\"algorithm_name\":\"行人闯入\",\"snowflake_id\":\"demo-001\"}" +``` + +## 7. 返回结构说明 + +接口始终返回统一结构: + +```json +{ + "code": 0, + "message": "ok", + "data": { + "event_id": "1768287987212271616", + "sms_preview": "发现有行人闯入行为,请去摄像头查看。", + "sms_ok": true, + "sms_skipped": false, + "receiver_count": 1, + "sms": { + "code": "OK", + "message": "OK", + "request_id": "...", + "biz_id": "..." + }, + "sms_error": "" + } +} +``` + +字段含义: +1. sms_preview:短信文本语义预览。 +2. sms_skipped:是否因为 LOG_ONLY_NO_SMS=true 而跳过实发。 +3. sms_ok:阿里云返回码是否为 OK。 +4. sms:阿里云网关返回详情。 +5. sms_error:发送失败时的错误信息。 + +## 8. 日志与排障 + +日志位置:logs/callback.log。 + +常见日志事件: +1. callback_received:收到回调。 +2. sms_preview:短信预览与模板参数。 +3. sms_receivers_loaded:当前收件人数。 +4. sms_sent:已调用阿里云并返回结果。 +5. sms_skipped:当前为仅预览模式。 +6. sms_config_missing:缺少必要配置。 + +常见问题: +1. 看见中文乱码:多数是终端显示编码问题,先以接口 JSON 返回为准。 +2. 未发送短信:检查 LOG_ONLY_NO_SMS 是否为 true。 +3. receiver_count 为 0:检查 SQLite 是否有 is_enabled=1 的号码。 +4. 阿里云报模板错误:确认签名、模板编码、模板变量 ACTION 与控制台一致。 + +## 9. 安全建议 + +1. 不要把 AK/SK 提交到版本库。 +2. 生产环境建议通过系统环境变量注入密钥。 +3. 定期轮换 AK/SK,并限制 RAM 权限最小化。 diff --git a/app/__pycache__/main.cpython-311.pyc b/app/__pycache__/main.cpython-311.pyc new file mode 100644 index 0000000..9767708 Binary files /dev/null and b/app/__pycache__/main.cpython-311.pyc differ diff --git a/app/main.py b/app/main.py new file mode 100644 index 0000000..9ee458c --- /dev/null +++ b/app/main.py @@ -0,0 +1,368 @@ +import json +import logging +import os +import sqlite3 +from logging.handlers import RotatingFileHandler +from typing import Any, Dict, List + +from fastapi import FastAPI, Request +from fastapi.responses import JSONResponse +from dotenv import load_dotenv + +from alibabacloud_tea_openapi import models as open_api_models +from alibabacloud_dysmsapi20170525.client import Client as Dysmsapi20170525Client +from alibabacloud_dysmsapi20170525 import models as dysmsapi_20170525_models + + +load_dotenv(override=True) + + +# Aliyun SMS config (prefer environment variables for sensitive values) +ALIYUN_ACCESS_KEY_ID = os.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID", "") +ALIYUN_ACCESS_KEY_SECRET = os.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET", "") +ALIYUN_REGION = os.getenv("ALIBABA_CLOUD_REGION_ID", "cn-hangzhou") +ALIYUN_SMS_SIGN_NAME = os.getenv("ALIYUN_SMS_SIGN_NAME", "") +ALIYUN_SMS_TEMPLATE_CODE = os.getenv("ALIYUN_SMS_TEMPLATE_CODE", "SMS_508915141") +SMS_DB_PATH = os.getenv("SMS_DB_PATH", "data/sms_receivers.db") +SMS_DB_TABLE = os.getenv("SMS_DB_TABLE", "sms_receivers") +# True: only log preview text, do not call Aliyun SMS API. +LOG_ONLY_NO_SMS = os.getenv("LOG_ONLY_NO_SMS", "false").strip().lower() in {"1", "true", "yes", "on"} + +SENSITIVE_FIELD_KEYS = { + "alarm_pic_data", + "src_pic_data", + "access_key_id", + "access_key_secret", + "phone", + "phone_numbers", + "mobile", +} + + +app = FastAPI(title="Callback to SMS") + + +def setup_logging() -> logging.Logger: + os.makedirs("logs", exist_ok=True) + + logger = logging.getLogger("callback_service") + logger.setLevel(logging.INFO) + + if logger.handlers: + return logger + + file_handler = RotatingFileHandler( + filename="logs/callback.log", + maxBytes=5 * 1024 * 1024, + backupCount=5, + encoding="utf-8", + ) + formatter = logging.Formatter( + "%(asctime)s | %(levelname)s | %(message)s", + datefmt="%Y-%m-%d %H:%M:%S", + ) + file_handler.setFormatter(formatter) + logger.addHandler(file_handler) + + console_handler = logging.StreamHandler() + console_handler.setFormatter(formatter) + logger.addHandler(console_handler) + + return logger + + +logger = setup_logging() + + +def parse_request_json(raw_body: bytes, content_type: str) -> Dict[str, Any]: + if not raw_body: + raise ValueError("empty body") + + encodings: List[str] = [] + content_type_lower = (content_type or "").lower() + if "charset=" in content_type_lower: + charset = content_type_lower.split("charset=", 1)[1].split(";", 1)[0].strip() + if charset: + encodings.append(charset) + + encodings.extend(["utf-8", "utf-8-sig", "gb18030", "gbk"]) + + tried = set() + for encoding in encodings: + if encoding in tried: + continue + tried.add(encoding) + try: + decoded_text = raw_body.decode(encoding) + parsed = json.loads(decoded_text) + if isinstance(parsed, dict): + return parsed + raise ValueError("json root must be object") + except Exception: + continue + + raise ValueError("invalid json or unsupported encoding") + + +def ensure_receivers_table() -> None: + os.makedirs(os.path.dirname(SMS_DB_PATH) or ".", exist_ok=True) + with sqlite3.connect(SMS_DB_PATH) as conn: + conn.execute( + f""" + CREATE TABLE IF NOT EXISTS {SMS_DB_TABLE} ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + phone_number TEXT NOT NULL UNIQUE, + is_enabled INTEGER NOT NULL DEFAULT 1, + updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP + ) + """ + ) + conn.commit() + + +def get_db_receivers() -> List[str]: + try: + with sqlite3.connect(SMS_DB_PATH) as conn: + rows = conn.execute( + f"SELECT phone_number FROM {SMS_DB_TABLE} WHERE is_enabled = 1 ORDER BY id" + ).fetchall() + return [str(row[0]).strip() for row in rows if str(row[0]).strip()] + except Exception: + logger.exception("sms_db_read_failed | db_path=%s | table=%s", SMS_DB_PATH, SMS_DB_TABLE) + return [] + + +def get_sms_receivers() -> List[str]: + db_receivers = get_db_receivers() + merged: List[str] = [] + seen = set() + for phone in db_receivers: + if phone in seen: + continue + seen.add(phone) + merged.append(phone) + return merged + + +def create_sms_client() -> Dysmsapi20170525Client: + config = open_api_models.Config( + access_key_id=ALIYUN_ACCESS_KEY_ID, + access_key_secret=ALIYUN_ACCESS_KEY_SECRET, + ) + # Dysmsapi endpoint is a fixed domain; region is carried by SDK config. + config.endpoint = "dysmsapi.aliyuncs.com" + config.region_id = ALIYUN_REGION + return Dysmsapi20170525Client(config) + + +def get_missing_sms_config(receivers: List[str]) -> List[str]: + missing: List[str] = [] + if not ALIYUN_ACCESS_KEY_ID: + missing.append("ALIBABA_CLOUD_ACCESS_KEY_ID") + if not ALIYUN_ACCESS_KEY_SECRET: + missing.append("ALIBABA_CLOUD_ACCESS_KEY_SECRET") + if not ALIYUN_SMS_SIGN_NAME: + missing.append("ALIYUN_SMS_SIGN_NAME") + if not ALIYUN_SMS_TEMPLATE_CODE: + missing.append("ALIYUN_SMS_TEMPLATE_CODE") + if not receivers: + missing.append("SQLite receivers table (enabled phone_number)") + return missing + + +def sanitize_payload(payload: Dict[str, Any]) -> Dict[str, Any]: + sanitized: Dict[str, Any] = {} + for key, value in payload.items(): + key_lower = str(key).lower() + if key_lower in SENSITIVE_FIELD_KEYS: + sanitized[key] = "***" + continue + if isinstance(value, str) and len(value) > 256: + sanitized[key] = f"" + continue + sanitized[key] = value + return sanitized + + +def send_sms_hardcoded(template_params: Dict[str, Any], phone_numbers: str) -> Dict[str, Any]: + client = create_sms_client() + request = dysmsapi_20170525_models.SendSmsRequest( + phone_numbers=phone_numbers, + sign_name=ALIYUN_SMS_SIGN_NAME, + template_code=ALIYUN_SMS_TEMPLATE_CODE, + # Keep escaped unicode to avoid transport/terminal charset side effects. + template_param=json.dumps(template_params, ensure_ascii=True), + ) + response = client.send_sms(request) + return { + "code": response.body.code, + "message": response.body.message, + "request_id": response.body.request_id, + "biz_id": response.body.biz_id, + } + + +def build_template_params(payload: Dict[str, Any]) -> Dict[str, Any]: + algorithm_name = str(payload.get("algorithm_name") or "").strip() + if not algorithm_name: + first_result = payload.get("result_data") + if isinstance(first_result, list) and first_result: + first_item = first_result[0] + if isinstance(first_item, dict): + algorithm_name = str(first_item.get("algorithm_name") or "").strip() + algorithm_name_en = str(payload.get("algorithm_name_en") or "Unknown").strip() or "Unknown" + + action_text = pick_safe_action_text(algorithm_name, algorithm_name_en) + + # These keys must match your approved SMS template variables. + return { + "ACTION": action_text, + } + + +def contains_suspicious_text(text: str) -> bool: + if not text: + return False + if "�" in text: + return True + # Private-use glyphs usually indicate decode issues in this scenario. + if any("\ue000" <= ch <= "\uf8ff" for ch in text): + return True + # Typical mojibake markers from UTF-8/Latin-1 mis-decoding. + if any(marker in text for marker in ("Ã", "Â", "æ", "ç", "å", "ä", "ï")): + cjk_count = sum(1 for ch in text if "\u4e00" <= ch <= "\u9fff") + if cjk_count == 0: + return True + return False + + +def pick_safe_action_text(algorithm_name: str, algorithm_name_en: str) -> str: + primary = (algorithm_name or "").strip() + fallback = (algorithm_name_en or "Unknown").strip() or "Unknown" + + if not primary: + return fallback + + # Only attempt repair when text looks suspicious; avoid changing healthy input. + if contains_suspicious_text(primary): + repaired = repair_mojibake_text(primary) + if repaired and not contains_suspicious_text(repaired): + logger.info("action_text_repaired | from=%s | to=%s", primary, repaired) + return repaired + + logger.warning("action_text_fallback | reason=suspicious_after_repair | original=%s | fallback=%s", primary, fallback) + return fallback + + return primary + + +def repair_mojibake_text(text: str) -> str: + """Repair common case: UTF-8 text incorrectly decoded as GBK.""" + if not text: + return text + for source_encoding in ("gb18030", "gbk", "cp936"): + try: + repaired = text.encode(source_encoding).decode("utf-8") + except Exception: + continue + + # Keep repaired value only when it differs and remains printable. + if repaired and repaired != text and all(ch.isprintable() for ch in repaired): + return repaired + return text + + +def build_sms_preview_text(template_params: Dict[str, Any]) -> str: + action = str(template_params.get("ACTION", "Unknown")) + return f"发现有{action}行为,请去摄像头查看。" + + +@app.on_event("startup") +async def startup_event() -> None: + ensure_receivers_table() + logger.info("sms_db_ready | path=%s | table=%s", SMS_DB_PATH, SMS_DB_TABLE) + + +@app.get("/healthz") +async def healthz() -> Dict[str, Any]: + receivers = get_sms_receivers() + return { + "status": "ok", + "receiver_count": len(receivers), + "receivers": receivers, + } + + +@app.post("/callback") +async def callback(request: Request) -> JSONResponse: + try: + raw_body = await request.body() + payload = parse_request_json(raw_body, request.headers.get("content-type", "")) + except Exception: + logger.exception("Invalid JSON received") + return JSONResponse( + status_code=400, + content={"code": 400, "message": "invalid json"}, + ) + + event_id = str(payload.get("snowflake_id") or payload.get("analysis_job_id") or "unknown") + safe_payload = sanitize_payload(payload) + logger.info("callback_received | event_id=%s | payload=%s", event_id, json.dumps(safe_payload, ensure_ascii=False)) + + params = build_template_params(payload) + sms_preview_text = build_sms_preview_text(params) + logger.info( + "sms_preview | event_id=%s | text=%s | template_param=%s", + event_id, + sms_preview_text, + json.dumps(params, ensure_ascii=False), + ) + + sms_ok = False + sms_skipped = False + sms_result: Dict[str, Any] = {} + sms_error = "" + receivers = get_sms_receivers() + receiver_count = len(receivers) + logger.info("sms_receivers_loaded | event_id=%s | count=%s", event_id, receiver_count) + + if LOG_ONLY_NO_SMS: + sms_skipped = True + logger.info("sms_skipped | event_id=%s | reason=log_only_mode", event_id) + else: + missing_config = get_missing_sms_config(receivers) + if missing_config: + sms_error = f"missing sms config: {', '.join(missing_config)}" + logger.error("sms_config_missing | event_id=%s | missing=%s", event_id, ",".join(missing_config)) + else: + try: + sms_result = send_sms_hardcoded(params, ",".join(receivers)) + sms_ok = str(sms_result.get("code", "")).upper() == "OK" + logger.info( + "sms_sent | event_id=%s | receiver_count=%s | sms_code=%s | sms_message=%s | biz_id=%s", + event_id, + receiver_count, + sms_result.get("code"), + sms_result.get("message"), + sms_result.get("biz_id"), + ) + except Exception as exc: + sms_error = str(exc) + logger.exception("sms_send_failed | event_id=%s | error=%s", event_id, sms_error) + + return JSONResponse( + status_code=200, + content={ + "code": 0, + "message": "ok", + "data": { + "event_id": event_id, + "sms_preview": sms_preview_text, + "sms_ok": sms_ok, + "sms_skipped": sms_skipped, + "receiver_count": receiver_count, + "sms": sms_result, + "sms_error": sms_error, + }, + }, + ) diff --git a/data/sms_receivers.db b/data/sms_receivers.db new file mode 100644 index 0000000..ad798cf Binary files /dev/null and b/data/sms_receivers.db differ diff --git a/logs/callback.log b/logs/callback.log new file mode 100644 index 0000000..b724e82 --- /dev/null +++ b/logs/callback.log @@ -0,0 +1 @@ +2026-07-07 10:13:56 | INFO | sms_db_ready | path=data/sms_receivers.db | table=sms_receivers diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..1a47e9a --- /dev/null +++ b/requirements.txt @@ -0,0 +1,5 @@ +fastapi==0.116.1 +uvicorn==0.35.0 +alibabacloud-dysmsapi20170525==3.1.1 +alibabacloud-tea-openapi==0.4.1 +python-dotenv==1.0.1