Webhook 自动化触发
Webhook 是「事件推送」:GitLab 合并、Harbor 镜像推送、Alertmanager 告警等 POST 到你的服务,由 Python 自动执行重启、同步、建单、通知等动作。本章实现一个可扩展的 Webhook 接收端。
整体流程
GitLab / Harbor / Alertmanager
│ POST /webhook/{source}
▼
FastAPI 路由(校验签名 / Token)
│
▼
校验 payload → 去重入库 → 投递持久队列
│
▼
Worker 调用脚本 / SSH / K8s API / 发企业微信
GitLab Push 示例
GitLab 项目 → Settings → Webhooks,URL 填 https://ops-api.internal/webhook/gitlab,Secret 与 GITLAB_WEBHOOK_SECRET 一致。生产任务不能依赖 BackgroundTasks:它在同一 Web 进程中运行,重启或扩缩容会丢失已接受任务。以下使用 Redis + RQ 演示“去重后入队”。
# API 与 worker 都需要相同的 Redis 地址和 Python 依赖。
uv add redis rq
export REDIS_URL="redis://redis.internal:6379/0"
#!/usr/bin/env python3
import logging
import os
import secrets
from fastapi import FastAPI, Header, HTTPException, Request, status
from redis import Redis
from rq import Queue
app = FastAPI()
logger = logging.getLogger("webhook")
GITLAB_SECRET = os.environ.get("GITLAB_WEBHOOK_SECRET", "")
redis = Redis.from_url(os.environ["REDIS_URL"], decode_responses=True)
queue = Queue("ops", connection=redis, default_timeout=900)
def verify_gitlab_token(x_gitlab_token: str | None) -> None:
"""GitLab 可在 Header 发送 X-Gitlab-Token(与 Secret 相同)。"""
if not GITLAB_SECRET or not x_gitlab_token:
raise HTTPException(status_code=403, detail="Webhook 校验失败")
if not secrets.compare_digest(x_gitlab_token, GITLAB_SECRET):
raise HTTPException(status_code=403, detail="Webhook 校验失败")
def enqueue_deploy(event_id: str, project: str, ref: str) -> bool:
"""以 Redis SET NX 做去重,成功后把任务交给独立 worker。"""
dedupe_key = f"webhook:seen:{event_id}"
if not redis.set(dedupe_key, "1", nx=True, ex=24 * 3600):
return False
try:
# Worker 函数必须位于可导入模块,例如 ops_worker.py。
queue.enqueue("ops_worker.run_deploy_pipeline", project, ref)
except Exception:
# 投递失败不能吞掉事件;移除标记以允许上游重试再次投递。
redis.delete(dedupe_key)
raise
return True
# ops_worker.py 中的函数:worker 进程负责真正的耗时变更。
def run_deploy_pipeline(project: str, ref: str) -> None:
"""实际实现应使用参数白名单、审计 ID 和受限执行账号。"""
logger.info("触发部署 project=%s ref=%s", project, ref)
# subprocess.run(["ansible-playbook", "deploy.yml", "-e", f"ref={ref}"], ...)
@app.post("/webhook/gitlab", status_code=status.HTTP_202_ACCEPTED)
async def gitlab_webhook(
request: Request,
x_gitlab_event: str | None = Header(default=None),
x_gitlab_token: str | None = Header(default=None),
):
verify_gitlab_token(x_gitlab_token)
body = await request.json()
# 只处理 push 到 main 的事件
if x_gitlab_event != "Push Hook":
return {"ignored": True, "reason": "非 Push 事件"}
ref = body.get("ref", "")
project = body.get("project", {}).get("path_with_namespace", "unknown")
if ref != "refs/heads/main":
return {"ignored": True, "reason": f"非 main 分支: {ref}"}
sha = body.get("checkout_sha", "")
if not sha:
raise HTTPException(status_code=422, detail="缺少 checkout_sha,无法安全去重")
event_id = f"gitlab:{project}:{ref}:{sha}"
accepted = enqueue_deploy(event_id, project, ref)
return {
"accepted": accepted,
"duplicate": not accepted,
"project": project,
"event_id": event_id,
}
要点: 校验并成功投递后返回 202 Accepted;worker 的重试、超时、死信队列和作业状态必须独立于 API 进程。BackgroundTasks 仅适合可丢失的短通知,不适合部署、删除和同步操作。
HMAC 签名校验(GitHub / 部分 Harbor)
import hashlib
import hmac
async def verify_hmac_sha256(request: Request, signature: str | None, secret: str) -> bytes:
"""Header 形如 sha256=abcdef..."""
raw = await request.body()
if not signature or not signature.startswith("sha256="):
raise HTTPException(status_code=403, detail="缺少签名")
expected = hmac.new(secret.encode(), raw, hashlib.sha256).hexdigest()
if not hmac.compare_digest(signature[7:], expected):
raise HTTPException(status_code=403, detail="签名不匹配")
return raw
Alertmanager 告警示例
Alertmanager receiver 配置 webhook URL,body 为 JSON 对象,单条告警位于顶层 alerts 数组。
@app.post("/webhook/alertmanager", status_code=status.HTTP_202_ACCEPTED)
async def alertmanager_webhook(request: Request):
# Alertmanager 的请求体是对象,单条告警位于 payload["alerts"],不是顶层数组。
payload = await request.json()
alerts = payload.get("alerts") if isinstance(payload, dict) else None
if not isinstance(alerts, list):
raise HTTPException(status_code=422, detail="Alertmanager payload 缺少 alerts 数组")
firing = [a for a in alerts if a.get("status") == "firing"]
for alert in firing:
name = alert.get("labels", {}).get("alertname", "unknown")
fingerprint = alert.get("fingerprint")
if fingerprint:
enqueue_alert(fingerprint, name, alert.get("annotations", {}))
return {"received": len(alerts), "firing": len(firing)}
def enqueue_alert(fingerprint: str, alertname: str, annotations: dict) -> None:
"""同一告警 fingerprint 在短窗口内只投递一次通知任务。"""
if redis.set(f"alert:seen:{fingerprint}", "1", nx=True, ex=300):
queue.enqueue("ops_worker.notify_oncall", alertname, annotations)
# ops_worker.py 中的函数:不要在 API 进程内直接发送耗时通知。
def notify_oncall(alertname: str, annotations: dict) -> None:
logger.warning("告警: %s %s", alertname, annotations.get("summary", ""))
# 对接钉钉 / 企业微信 / PagerDuty
幂等与防重放
| 策略 | 说明 |
|---|---|
| 事件 ID 去重 | GitLab 项目、分支和 checkout_sha 组合写入 Redis,TTL 24h |
| 快速响应 | 完成校验与入队后返回 202,任务进队列(Celery / RQ) |
| 失败重试 | 由 worker 重试;队列应配置最大次数、死信与告警 |
| 日志 | 保留原始 payload 摘要(勿记录 Secret) |
# 简易内存去重(单进程够用;多副本用 Redis)
_seen: set[str] = set()
def dedupe(event_id: str) -> bool:
if event_id in _seen:
return True
_seen.add(event_id)
return False
提示
第八阶段综合项目会把 Webhook、认证、批量任务组合成完整平台;本章接口可直接作为 api/webhook.py 模块起点。