AbhiramiRajeev/Ingestion-Service

GitHub: AbhiramiRajeev/Ingestion-Service

一个基于 Go 和 Gin 的安全事件接入服务,通过 HTTP 收集并验证事件数据后将其发布到 Kafka 以供下游处理。

Stars: 0 | Forks: 0

# Ingestion 服务 ## 问题陈述 现代安全系统会实时地从多个端点产生大量事件 —— 登录尝试、可疑活动、暴力破解信号等。在不丢失数据或产生瓶颈的情况下,收集、验证并可靠地将这些事件转发到处理流水线至关重要。 该服务通过提供一个轻量级、经过身份验证的接入层来解决此问题,该层通过 HTTP 接收事件,对其进行验证,并立即将其发布到 Kafka,供下游消费者进行处理。 ## 架构 ``` ┌──────────────┐ │ Client / │ │ Agent │ └──────┬───────┘ │ POST /ingest │ Authorization: Bearer ▼ ┌──────────────┐ │ Ingestion │ │ Service │ │ (Gin) │ └──────┬───────┘ │ Validate API Key │ Bind + Validate Payload │ Marshal to JSON ▼ ┌──────────────┐ │ Sync Kafka │ │ Producer │ └──────┬───────┘ │ Publish to topic ▼ ┌──────────────┐ │ Kafka │ │ Broker │ │ (login_events)│ └──────────────┘ ``` ## 请求流程 ``` API Request │ ▼ JWT Validation │ ▼ Validate Payload │ ▼ Publish to Kafka │ ▼ Return 200 OK ``` ## API 端点 ### `POST /ingest` 接收安全事件并将其发布到 Kafka。 **请求头:** | Header | 值 | 必填 | |---------------|------------------------|----------| | Authorization | `Bearer ` | 是 | | Content-Type | `application/json` | 是 | **请求体:** ``` { "event_type": "login_attempt", "user_id": "alice", "ip_address": "1.2.3.4", "status": "failure", "timestamp": "2025-06-16T12:00:00Z" } ``` **响应 (200 OK):** ``` { "message": "Event ingested successfully" } ``` **错误响应:** | 状态码 | Body | 何时触发 | |--------|-----------------------------------------|----------------------------| | 401 | `{ "error": "Unauthorized" }` | API key 缺失或无效 | | 400 | `{ "error": "Invalid request body" }` | 格式错误的 JSON 或缺少必填字段 | | 500 | `{ "error": "Failed to send message to Kafka" }` | Kafka 发布失败 | ### `GET /health` 健康检查端点。目前为一个存根。 ## 身份验证流程 ``` Client sends request │ ▼ Extract "Authorization" header │ ▼ Split header into "Bearer" + │ ▼ Compare token against configured API_KEY │ ├── Match → proceed with request │ └── No match → return 401 Unauthorized ``` 身份验证是一个简单的 API key 检查 —— 客户端必须在 `Authorization` 请求头中发送有效的 `API_KEY`(设置为环境变量)。这里没有 JWT 解码;而是直接将原始 token 与该 key 进行比对。 ## Kafka 事件格式 事件以 JSON 格式发布到 `login_events` 主题,其 schema 如下: | 字段 | 类型 | 描述 | |--------------|--------|-------------------------------------------| | `event_type` | string | 事件类型(例如 `login_attempt`) | | `user_id` | string | 相关用户的标识符 | | `ip_address` | string | 事件的源 IP 地址 | | `status` | string | 结果(例如 `success`, `failure`) | | `timestamp` | string | 事件的 ISO 8601 时间戳 | **Kafka 消息示例:** ``` { "event_type": "login_attempt", "user_id": "alice", "ip_address": "1.2.3.4", "status": "failure", "timestamp": "2025-06-16T12:00:00Z" } ``` ## 目录结构 ``` . ├── cmd/ │ └── main.go # Entrypoint — wires Kafka, Gin router, and handlers ├── config/ │ ├── config.go # Config struct definition │ └── config.yaml # Default configuration values ├── internal/ │ ├── auth.go # API key validation logic │ ├── handler.go # HTTP handlers (IngestEvent, GetHealth) │ └── kafka.go # Kafka producer initialization ├── models/ │ └── models.go # Event data model ├── Dockerfile # Multi-stage build for production image ├── docker-compose.yml # Zookeeper + Kafka for local dev ├── go.mod └── go.sum ``` ## 配置 配置通过 `config/config.yaml` 加载: ``` kafka: brokers: - "localhost:9092" topic: - "login_events" auth: api_key: - "my-api-key" server: port: - "8080" ``` | 设置项 | 描述 | 默认值 | |-----------------|------------------------------------------|-----------------| | `kafka.brokers` | Kafka broker 地址列表 | `localhost:9092`| | `kafka.topic` | 用于发布事件的 Kafka 主题 | `login_events` | | `auth.api_key` | 用于身份验证的允许 API key | — | | `server.port` | HTTP 服务器监听的端口 | `8080` | ## Docker ### 在本地构建并运行 ``` # 构建 image docker build -t ingestion-service . # 运行 container docker run -p 8080:8080 -e API_KEY=my-secret-key ingestion-service ``` ### 使用 docker-compose (Kafka + Zookeeper) ``` docker-compose up -d ``` 这将启动: | 服务 | 端口 | |-------------|------| | Zookeeper | 2181 | | Kafka | 9092 | ### 完整的本地设置 ``` # 1. 启动 Kafka docker-compose up -d # 2. 设置 API key export API_KEY=my-secret-key # 3. 运行 service go run cmd/main.go ``` ## 技术栈 | 组件 | 技术 | |---------------|-----------------------------| | HTTP 服务器 | [Gin](https://github.com/gin-gonic/gin) | | Kafka 客户端 | [Sarama](https://github.com/IBM/sarama) | | 日志 | [klog](https://github.com/kubernetes/klog) | | 语言 | Go 1.24 | ## 未来改进 - **JWT 验证** — 将简单的 API key 检查替换为 JWT token 验证和基于声明的授权 - **异步 producer** — 切换到异步 Kafka producer 以在负载下获得更高的吞吐量 - **结构化日志** — 在日志中添加请求 ID 和结构化字段 - **指标** — 暴露 Prometheus 指标(请求计数、延迟、Kafka 发布速率) - **限流** — 添加基于客户端的限流以防止滥用 - **Schema 验证** — 根据 JSON schema 或 protobuf 定义验证事件 payload - **死信队列 (DLQ)** — 将失败的消息路由到 DLQ 以便稍后重试 - **Kafka consumer** — 构建一个 consumer 服务来处理来自 `login_events` 的事件 - **优雅关闭** — 处理 SIGTERM 以清理进行中的请求并干净地关闭 Kafka - **集成测试** — 添加带有测试 Kafka broker 的测试覆盖率
标签:API网关, EVTX分析, Gin, Go, Kafka, Ruby工具, SonarQube插件, 安全事件采集, 安全编排与自动化响应, 日志审计, 请求拦截