mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-28 05:46:36 +08:00
Compare commits
10 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| f960511cc0 | |||
| 465440fa5b | |||
| a4dd5ca9e1 | |||
| a9e4237bbf | |||
| 75d1fcf345 | |||
| f1577bf092 | |||
| 01ed2c5e36 | |||
| 80696c12fa | |||
| 3d4d99081e | |||
| 0639855653 |
@@ -28,17 +28,14 @@ description: "Wavelet 项目专用:根据自上一个正式版本 Tag 以来
|
||||
|
||||
1. 合并重复或相近提交。
|
||||
2. 删除无意义提交,例如格式化、临时调试、无关重构。
|
||||
3. 将内部实现描述改写为用户可理解的变更。
|
||||
4. 每条使用完整中文句子。
|
||||
5. 尽量说明“修复/优化了什么”以及“带来的效果”。
|
||||
6. 不要编造 commit log 中没有的信息。
|
||||
7. 不要加入 token、密钥、私有地址等敏感信息。
|
||||
8. 如果某个分类没有内容,可以省略。
|
||||
3. 将内部实现描述改写为用户可理解的变更, 说明“修复/优化了什么”以及“带来的效果”。
|
||||
4. 不要写技术细节:只描述用户可感知的行为与效果,禁止内部实现描述,例如字段名/表名/SQL(`node_id = ''`)、框架或库名称(shadcn、GORM、OpenResty)、配置或协议细节(RFC3339、ClickHouse/PostgreSQL 差异)、代码机制(`proxy_intercept_errors`、Lua 过滤器、雪花 ID)。数据库名称仅在说明受影响用户范围时使用(如「PostgreSQL 日志库下无数据」)。
|
||||
5. 如果某个分类没有内容,则省略。
|
||||
|
||||
固定使用以下分类:
|
||||
|
||||
```text
|
||||
### 新增
|
||||
### ✨ 新功能
|
||||
### 🛠 修复
|
||||
### ⚡️ 优化与改进
|
||||
### 💄 其他/体验
|
||||
@@ -46,7 +43,7 @@ description: "Wavelet 项目专用:根据自上一个正式版本 Tag 以来
|
||||
|
||||
分类规则:
|
||||
|
||||
- 新功能、新能力、新配置、新任务:放入 ### 新增
|
||||
- 新功能、新能力、新配置、新任务:放入 ### ✨ 新功能
|
||||
- Bug、异常行为、错误逻辑:放入 ### 🛠 修复
|
||||
- 性能、稳定性、接口、架构、兼容性:放入 ### ⚡️ 优化与改进
|
||||
- 日志、文案、UI、文档、开发体验:放入 ### 💄 其他/体验
|
||||
@@ -67,9 +64,8 @@ chore(release): v3.3.0
|
||||
- 新增笔记库快照备份功能,支持定时备份与手动一键恢复(仅说明新增的功能, 禁止提及新功能开发时期的优化修复等内容)。
|
||||
|
||||
### 🛠 修复
|
||||
- 修复了通过 MCP 接口操作时笔记库范围限制未正确生效的问题。
|
||||
- 修复了 MCP 接口返回数据格式不一致的问题。
|
||||
- 修复了 WebSocket 客户端异常断开后僵尸连接未及时清理的问题。
|
||||
- 修复首页「来源分布」卡片在 PostgreSQL/SQLite 日志库下无数据的问题。
|
||||
- 修复源站错误页「仅针对 GET 请求」未真正透传非 GET 响应的问题:POST/PUT 等非 GET 请求现可完整看到源站原始报错内容。
|
||||
|
||||
### ⚡️ 优化与改进
|
||||
- 优化了 WebGUI 登录机制,引入设备令牌自动轮转,减少因 IP 变化产生的冗余令牌。
|
||||
|
||||
@@ -1,19 +0,0 @@
|
||||
root = true
|
||||
|
||||
[*]
|
||||
indent_style = space
|
||||
indent_size = 4
|
||||
charset = utf-8
|
||||
end_of_line = lf
|
||||
trim_trailing_whitespace = true
|
||||
insert_final_newline = true
|
||||
|
||||
[*.{json,yml,yaml}]
|
||||
indent_size = 2
|
||||
|
||||
[*.md]
|
||||
insert_final_newline = false
|
||||
trim_trailing_whitespace = false
|
||||
|
||||
[*.{js,ts,css,html,jsx,tsx,vue}]
|
||||
indent_size = 2
|
||||
+61
-106
@@ -2,9 +2,9 @@
|
||||
|
||||
# OpenFlare
|
||||
|
||||
**[English](./README.en.md) | [📖 中文](./README.md)**
|
||||
**[📖 中文](./README.md) | [English](./README.en.md)**
|
||||
|
||||
OpenFlare is an open-source CDN orchestration and edge security platform. It supports reverse proxies, centralized configuration synchronization, secure intranet penetration (Tunnels), dynamic WAF protection, and anti-CC challenges.
|
||||
OpenFlare is an open-source CDN orchestration and edge security platform. It supports reverse proxy, centralized configuration synchronization, in-network tunneling (Tunnels), dynamic WAF protection, and CC defense challenges.
|
||||
|
||||
</div>
|
||||
|
||||
@@ -21,37 +21,67 @@ OpenFlare is an open-source CDN orchestration and edge security platform. It sup
|
||||
</p>
|
||||
|
||||
> [!WARNING]
|
||||
> After logging in for the first time with the `root` user, make sure to change the default password `123456`.
|
||||
>
|
||||
> The BETA version is a temporary product for the development and testing phase. It may contain unknown issues and should not be used in production environments.
|
||||
> After the first login with the `admin` user, you must change the default password `12345678`.
|
||||
>
|
||||
> The BETA version is a temporary product in the development and testing stage and may have unknown issues. It should not be used in production environments.
|
||||
|
||||
## Documentation
|
||||
|
||||
**https://open-flare.pages.dev**
|
||||
|
||||
Quick links:
|
||||
Common entry points:
|
||||
|
||||
* [Quick Start](https://open-flare.pages.dev/en/guide/quick-start)
|
||||
* [Deployment Guide](https://open-flare.pages.dev/en/deployment/deployment)
|
||||
* [Quick Start](https://open-flare.pages.dev/guide/quick-start)
|
||||
* [Deployment Guide](https://open-flare.pages.dev/deployment/deployment)
|
||||
* [Configuration Reference](https://open-flare.pages.dev/reference/configuration)
|
||||
* [System Design](https://open-flare.pages.dev/design/)
|
||||
|
||||
## Core Features
|
||||
## Core Capabilities
|
||||
|
||||
* **Reverse Proxy Management**: Website rules as the aggregation boundary, supporting multi-domain binding and multi-upstream load balancing with unified management of all OpenResty node configurations.
|
||||
* **Immutable Config Version Control**: Full-snapshot publish model based on version numbers (`YYYYMMDD-NNN`), with pre-publish diff preview, a single globally active version, and one-click sub-second rollback.
|
||||
* **Secure Intranet Penetration (Tunnels)**: An open-source alternative to Cloudflare Tunnels. Securely expose local intranet Web services to the public network via Relay and OpenFlared clients — no public IP or open inbound ports required.
|
||||
* **Edge WAF Safety Protection**: Provides global and custom rule groups, supporting manual/automatic/subscription IP groups, MaxMind GeoIP country-level access control, Checksum-based differential IP group sync (no Nginx reload), and custom block responses.
|
||||
* **Anti-CC & Human-Machine Challenge (PoW)**: Built-in high-performance client-side cryptographic Proof of Work challenges (similar to Turnstile) to block and intercept botnets and scrapers at the gateway edge in seconds.
|
||||
* **Pages Static Hosting**: Upload pre-built ZIP packages directly; edge Agents pull and serve them via local OpenResty, with SPA Fallback and built-in API reverse proxy configuration.
|
||||
* **Automated TLS Certificate Management**: Supports dynamic certificate upload, automatic multi-domain certificate matching and binding, and ACME-based automatic issuance and renewal via Let's Encrypt.
|
||||
* **Uptime Kuma Monitoring Sync**: Integrates with Uptime Kuma to automatically sync the monitoring site list using differential updates, providing real-time awareness of node availability and service health.
|
||||
* **SSO Single Sign-On**: Supports GitHub OAuth and standard OIDC protocol for seamless integration with enterprise identity providers.
|
||||
* **Unified Observability**: Aggregates node request metrics, real-time access log details, host/Nginx resource snapshots, health events, and a re-upload buffer for network fluctuations.
|
||||
* **Reverse Proxy Configuration Management**: Uses website rules as the aggregation boundary, supports multi-domain binding and multi-upstream load balancing, and centrally manages reverse proxy configurations for all OpenResty nodes.
|
||||
* **Secure In-Network Tunneling (Tunnels)**: Open-source version of Cloudflare Tunnels. No public IP or exposed inbound ports are required. Securely reverse-proxy internal web services to the public internet through Relay relay nodes and OpenFlared clients.
|
||||
* **Edge WAF Security Protection**: Provides global and custom rule groups, supports manual/auto/subscription-type IP groups, MaxMind GeoIP national-level geographic access control, IP group member Checksum differential synchronization (no Nginx reload required), and custom blocking responses.
|
||||
* **CC Defense and Human-Computer Challenge (PoW)**: Built-in high-performance client-side cryptography Proof of Work challenge (similar to Turnstile). Secures high-speed interception and blocking of zombie networks and crawlers at the gateway edge.
|
||||
* **Pages Static Hosting**: Supports uploading or synchronizing pre-built artifacts from restricted Remote URLs or public GitHub Release assets. GitHub latest can be checked periodically and optionally auto-published. All sources are unified to generate immutable deployments, pulled by the edge Agent and served locally by OpenResty, supporting rollbacks, SPA Fallback, and API reverse proxy.
|
||||
* **TLS Certificate Automation**: Supports dynamic certificate uploads, automatic multi-domain certificate matching and binding, and automatic issuance and renewal of certificates from Let's Encrypt via the ACME protocol.
|
||||
* **Uptime Kuma Monitoring Synchronization**: Integrated with Uptime Kuma to automatically perform differential synchronization of monitoring site lists, real-time awareness of node availability and service status.
|
||||
* **SSO Single Sign-On**: Supports GitHub OAuth and standard OIDC protocol for seamless integration with enterprise identity providers to achieve unified login.
|
||||
* **Unified Observability**: Aggregates node request metrics, real-time access log details, host and Nginx resource snapshots, health events, and network fluctuation replenishment buffers.
|
||||
|
||||
## Interface Preview
|
||||
|
||||
### Dashboard Overview
|
||||
|
||||

|
||||
|
||||
### Access Logs
|
||||
|
||||

|
||||
|
||||
### WAF Protection
|
||||
|
||||

|
||||
|
||||
## Quick Start
|
||||
|
||||
### 1. Launch Server
|
||||
### Hardware Configuration Recommendations
|
||||
|
||||
| Component | Minimum Hardware Requirements | Recommended Hardware Requirements | Notes |
|
||||
|------------------------|-----------------------------------|-----------------------------------|-------|
|
||||
| **Server Control Plane** | 1 CPU core / 2 GB RAM / 20 GB disk | 2 CPU cores / 4 GB RAM / 50 GB+ disk | Disk usage should be expanded reasonably based on access log retention duration and concurrent traffic |
|
||||
| **Agent Data Plane** | 1 CPU core / 512 MB RAM / 2 GB disk | 2 CPU cores / 2 GB RAM / 10 GB+ disk | Expanded based on OpenResty concurrent proxy connections and WAF interception processing |
|
||||
| **Relay Relay Node** | 1 CPU core / 1 GB RAM / 5 GB disk | 2 CPU cores / 2 GB RAM / 20 GB disk | frps transmission relay throughput is mainly limited by bandwidth and CPU throughput |
|
||||
| **OpenFlared Client** | 1 CPU core / 256 MB RAM / 1 GB disk | 1 CPU core / 512 MB RAM / 5 GB disk | Runs independently on the internal network with extremely low resource consumption; only network throughput needs to be guaranteed |
|
||||
|
||||
### 1. Start the Server
|
||||
|
||||
Use `docker-compose`:
|
||||
|
||||
```bash
|
||||
# Download environment variable template and create .env file
|
||||
curl -o .env.example https://raw.githubusercontent.com/Rain-kl/OpenFlare/refs/heads/main/.env.example
|
||||
cp .env.example .env
|
||||
```
|
||||
|
||||
```yaml
|
||||
services:
|
||||
@@ -70,8 +100,6 @@ services:
|
||||
condition: service_healthy
|
||||
redis:
|
||||
condition: service_healthy
|
||||
clickhouse:
|
||||
condition: service_healthy
|
||||
|
||||
postgres:
|
||||
image: postgres:17-alpine
|
||||
@@ -101,118 +129,45 @@ services:
|
||||
retries: 5
|
||||
start_period: 5s
|
||||
|
||||
clickhouse:
|
||||
image: clickhouse/clickhouse-server:25.3-alpine
|
||||
restart: unless-stopped
|
||||
environment:
|
||||
CLICKHOUSE_DB: ${CLICKHOUSE_NAME:-openflare}
|
||||
CLICKHOUSE_USER: ${CLICKHOUSE_USERNAME:-default}
|
||||
CLICKHOUSE_PASSWORD: ${CLICKHOUSE_PASSWORD:-replace-with-clickhouse-password}
|
||||
CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: 1
|
||||
TZ: ${TZ:-Asia/Shanghai}
|
||||
volumes:
|
||||
- openflare_clickhouse_data:/var/lib/clickhouse
|
||||
healthcheck:
|
||||
test: ["CMD", "clickhouse-client", "--user", "${CLICKHOUSE_USERNAME:-default}", "--password", "${CLICKHOUSE_PASSWORD:-replace-with-clickhouse-password}", "--query", "SELECT 1"]
|
||||
interval: 10s
|
||||
timeout: 5s
|
||||
retries: 5
|
||||
start_period: 15s
|
||||
|
||||
volumes:
|
||||
openflare_uploads:
|
||||
openflare_postgres_data:
|
||||
openflare_redis_data:
|
||||
openflare_clickhouse_data:
|
||||
```
|
||||
|
||||
```bash
|
||||
docker compose up -d
|
||||
```
|
||||
See the [deployment documentation](https://open-flare.pages.dev/deployment/deployment) for details.
|
||||
|
||||
Access at: `http://localhost:3000`
|
||||
Access address: `http://localhost:3000`
|
||||
|
||||
Default credentials:
|
||||
Default account:
|
||||
|
||||
* Username: `root`
|
||||
* Password: `123456`
|
||||
* Username: `admin`
|
||||
* Password: `12345678`
|
||||
|
||||
### 2. Install Agent
|
||||
|
||||
Before installing an Agent, please install OpenResty on the target node first, or use the Agent Docker image with OpenResty built-in.
|
||||
Before installing the Agent, first install OpenResty on the node or use the built-in OpenResty Agent Docker image.
|
||||
|
||||
You can copy the installation command from **Node Management -> Details -> Node Info -> Node Token & Deployment** in the control panel, or directly use the scripts below:
|
||||
You can copy the installation command from the control panel's **Nodes Management -> Details -> Node Information -> Node ID and Deployment**, or use the script below:
|
||||
|
||||
#### Docker Deployment
|
||||
|
||||
For Docker deployment, you can directly run the Agent image:
|
||||
Docker deployment can directly run the Agent image:
|
||||
|
||||
```bash
|
||||
docker pull ghcr.io/rain-kl/openflare-agent:latest
|
||||
docker rm -f openflare-agent 2>/dev/null || true
|
||||
docker run -d --name openflare-agent --restart unless-stopped \
|
||||
-p 80:80 -p 443:443/tcp -p 443:443/udp \
|
||||
-v openflare-agent-pages:/data/var/lib/openflare/pages \
|
||||
-e OPENFLARE_SERVER_URL=http://your-server:3000 \
|
||||
-e OPENFLARE_AGENT_TOKEN=YOUR_AGENT_TOKEN \
|
||||
ghcr.io/rain-kl/openflare-agent:latest
|
||||
```
|
||||
|
||||
#### Local Installation
|
||||
## Open Source License
|
||||
|
||||
Using `discovery_token` to register:
|
||||
|
||||
```bash
|
||||
curl -fsSL https://raw.githubusercontent.com/Rain-kl/OpenFlare/main/scripts/install-agent.sh | bash -s -- \
|
||||
--server-url http://your-server:3000 \
|
||||
--discovery-token YOUR_DISCOVERY_TOKEN
|
||||
```
|
||||
|
||||
Using node-specific `agent_token`:
|
||||
|
||||
```bash
|
||||
curl -fsSL https://raw.githubusercontent.com/Rain-kl/OpenFlare/main/scripts/install-agent.sh | bash -s -- \
|
||||
--server-url http://your-server:3000 \
|
||||
--agent-token YOUR_AGENT_TOKEN
|
||||
```
|
||||
|
||||
The installation script defaults to `/opt/openflare-agent`, creates a `openflare-agent.service`, automatically searches for `openresty`, and can be executed repeatedly to reinstall or upgrade the Agent.
|
||||
|
||||
### 3. Uninstall Agent
|
||||
|
||||
To completely uninstall the Agent and clear local data, run:
|
||||
|
||||
```bash
|
||||
curl -fsSL https://raw.githubusercontent.com/Rain-kl/OpenFlare/main/scripts/uninstall-agent.sh | bash
|
||||
```
|
||||
|
||||
The uninstallation script will stop and remove the `openflare-agent.service`, and delete the entire `/opt/openflare-agent` directory. It will not delete the local OpenResty installation.
|
||||
|
||||
### 4. Publish Your First Configuration
|
||||
|
||||
1. Log in to the management panel and add a reverse proxy rule.
|
||||
2. View the preview or change summary before publishing.
|
||||
3. Activate the new version.
|
||||
4. Agents will receive the configuration and apply it via WebSocket notification or subsequent heartbeats.
|
||||
|
||||
The version number format is fixed as `YYYYMMDD-NNN`. Historical versions are immutable, and rollback is achieved by reactivating an older version.
|
||||
|
||||
## UI Preview
|
||||
|
||||
### Dashboard Overview
|
||||
|
||||

|
||||
|
||||
### Node Details
|
||||
|
||||

|
||||
|
||||
### Proxy Configuration
|
||||
|
||||

|
||||
|
||||
## License
|
||||
|
||||
This project is licensed under [Apache License 2.0](./LICENSE).
|
||||
This project is licensed under the [Apache License 2.0](./LICENSE).
|
||||
|
||||
## Star History
|
||||
|
||||
|
||||
+15
-1
@@ -16,7 +16,21 @@ sidebar: false
|
||||
>
|
||||
|
||||
|
||||
## [Unreleased]
|
||||
## [v3.5.3] - 2026-08-13
|
||||
|
||||
### 新增
|
||||
- 访问日志「日志明细」支持按 HTTP 状态码筛选,可直接输入任意状态码。
|
||||
- 访问日志「日志明细」支持自定义时间范围筛选,可按起止时间检索日志。
|
||||
- 首页看板改版:24 小时请求趋势拆分展示请求总量与 2xx/4xx/5xx 状态码类请求量并独占一行;移除宿主机磁盘指标,24 小时容量趋势(CPU/内存)并入业务流量卡片展示。
|
||||
|
||||
### 🛠 修复
|
||||
- 修复首页「来源分布」卡片在 PostgreSQL/SQLite 日志库下无数据的问题。
|
||||
- 修复源站错误页「仅针对 GET 请求」未真正透传非 GET 响应的问题:POST/PUT 等非 GET 请求现可完整看到源站原始报错内容。
|
||||
|
||||
## [v3.5.2] - 2026-08-09
|
||||
|
||||
### 🛠 修复
|
||||
- 修复 PostgreSQL 作为日志库时节点访问日志/可观测指标/用户访问日志批量写入失败的问题,现可正常写入。
|
||||
|
||||
## [v3.5.1] - 2026-08-09
|
||||
|
||||
|
||||
+19
-1
@@ -4809,7 +4809,7 @@ const docTemplate = `{
|
||||
"SessionCookie": []
|
||||
}
|
||||
],
|
||||
"description": "分页返回 OpenFlare 访问日志,支持按节点、IP、主机与路径筛选,需要管理员权限",
|
||||
"description": "分页返回 OpenFlare 访问日志,支持按节点、IP、主机、路径与状态码筛选,需要管理员权限",
|
||||
"produces": [
|
||||
"application/json"
|
||||
],
|
||||
@@ -4842,6 +4842,24 @@ const docTemplate = `{
|
||||
"name": "path",
|
||||
"in": "query"
|
||||
},
|
||||
{
|
||||
"type": "integer",
|
||||
"description": "HTTP 状态码(100-599)",
|
||||
"name": "status_code",
|
||||
"in": "query"
|
||||
},
|
||||
{
|
||||
"type": "string",
|
||||
"description": "起始时间(RFC3339,需与 until 成对提供)",
|
||||
"name": "since",
|
||||
"in": "query"
|
||||
},
|
||||
{
|
||||
"type": "string",
|
||||
"description": "结束时间(RFC3339,需与 since 成对提供)",
|
||||
"name": "until",
|
||||
"in": "query"
|
||||
},
|
||||
{
|
||||
"type": "integer",
|
||||
"description": "页码",
|
||||
|
||||
@@ -264,7 +264,7 @@ Server 的所有核心基础配置定义在 `config.yaml` 中,且均支持环
|
||||
| 配置键 (Key) | 数据类型 | 作用说明 | 默认值 |
|
||||
| --- | --- | --- | --- |
|
||||
| `origin_error_page_enabled` | `bool` | 是否启用全局源站错误页。开启后,源站或网关返回的匹配状态码由自定义/默认 HTML 替换,**HTTP 状态码保持原值**;关闭后不生成相关指令,恢复透传。修改后需发布配置版本生效 | `true` |
|
||||
| `origin_error_page_get_only` | `bool` | 是否仅对 **GET** 请求生效。开启后仅 GET 的匹配错误状态码返回自定义错误页;POST/PUT 等其它方法不返回自定义错误页(保留原始错误状态码) | `false` |
|
||||
| `origin_error_page_get_only` | `bool` | 是否仅对 **GET** 请求生效。开启后仅 GET 的匹配错误状态码返回自定义错误页;POST/PUT 等其它方法**透传源站响应**(原始状态码与响应体不变) | `false` |
|
||||
| `origin_error_page_status_codes` | `json` | 触发错误页的状态码标签 JSON 数组。支持单码(如 `522`)与闭区间(如 `500-599`);单码与区间两端均须在 **400–599**,且 `lo ≤ hi`。启用时展开结果不能为空 | `["500-599"]` |
|
||||
| `origin_error_page_html` | `string` | 错误页自定义 HTML。空字符串表示使用内置 OpenFlare 默认模板(极简白底);支持占位符 `{{status}}`(与 HTTP 状态码一致)、`{{host}}`(请求 Host)。最大 **256 KiB**(按字节)。勿嵌入不可信第三方脚本 | 空 |
|
||||
|
||||
|
||||
+19
-1
@@ -4802,7 +4802,7 @@
|
||||
"SessionCookie": []
|
||||
}
|
||||
],
|
||||
"description": "分页返回 OpenFlare 访问日志,支持按节点、IP、主机与路径筛选,需要管理员权限",
|
||||
"description": "分页返回 OpenFlare 访问日志,支持按节点、IP、主机、路径与状态码筛选,需要管理员权限",
|
||||
"produces": [
|
||||
"application/json"
|
||||
],
|
||||
@@ -4835,6 +4835,24 @@
|
||||
"name": "path",
|
||||
"in": "query"
|
||||
},
|
||||
{
|
||||
"type": "integer",
|
||||
"description": "HTTP 状态码(100-599)",
|
||||
"name": "status_code",
|
||||
"in": "query"
|
||||
},
|
||||
{
|
||||
"type": "string",
|
||||
"description": "起始时间(RFC3339,需与 until 成对提供)",
|
||||
"name": "since",
|
||||
"in": "query"
|
||||
},
|
||||
{
|
||||
"type": "string",
|
||||
"description": "结束时间(RFC3339,需与 since 成对提供)",
|
||||
"name": "until",
|
||||
"in": "query"
|
||||
},
|
||||
{
|
||||
"type": "integer",
|
||||
"description": "页码",
|
||||
|
||||
+13
-1
@@ -7194,7 +7194,7 @@ paths:
|
||||
- config
|
||||
/api/v1/d/access-logs:
|
||||
get:
|
||||
description: 分页返回 OpenFlare 访问日志,支持按节点、IP、主机与路径筛选,需要管理员权限
|
||||
description: 分页返回 OpenFlare 访问日志,支持按节点、IP、主机、路径与状态码筛选,需要管理员权限
|
||||
parameters:
|
||||
- description: 节点 ID
|
||||
in: query
|
||||
@@ -7212,6 +7212,18 @@ paths:
|
||||
in: query
|
||||
name: path
|
||||
type: string
|
||||
- description: HTTP 状态码(100-599)
|
||||
in: query
|
||||
name: status_code
|
||||
type: integer
|
||||
- description: 起始时间(RFC3339,需与 until 成对提供)
|
||||
in: query
|
||||
name: since
|
||||
type: string
|
||||
- description: 结束时间(RFC3339,需与 since 成对提供)
|
||||
in: query
|
||||
name: until
|
||||
type: string
|
||||
- description: 页码
|
||||
in: query
|
||||
name: p
|
||||
|
||||
@@ -1,9 +1,22 @@
|
||||
'use client';
|
||||
|
||||
import { Search } from 'lucide-react';
|
||||
import { useState } from 'react';
|
||||
import { format } from 'date-fns';
|
||||
import { CalendarIcon, ChevronDown, Search } from 'lucide-react';
|
||||
|
||||
import { Button } from '@/components/ui/button';
|
||||
import { Calendar } from '@/components/ui/calendar';
|
||||
import {
|
||||
Collapsible,
|
||||
CollapsibleContent,
|
||||
CollapsibleTrigger,
|
||||
} from '@/components/ui/collapsible';
|
||||
import { Input } from '@/components/ui/input';
|
||||
import {
|
||||
Popover,
|
||||
PopoverContent,
|
||||
PopoverTrigger,
|
||||
} from '@/components/ui/popover';
|
||||
import {
|
||||
Select,
|
||||
SelectContent,
|
||||
@@ -23,6 +36,115 @@ interface AccessLogFiltersProps {
|
||||
onReset: () => void;
|
||||
}
|
||||
|
||||
const HOUR_OPTIONS = Array.from({ length: 24 }, (_, i) =>
|
||||
String(i).padStart(2, '0'),
|
||||
);
|
||||
const MINUTE_OPTIONS = Array.from({ length: 60 }, (_, i) =>
|
||||
String(i).padStart(2, '0'),
|
||||
);
|
||||
|
||||
function FilterField({
|
||||
label,
|
||||
children,
|
||||
}: {
|
||||
label: string;
|
||||
children: React.ReactNode;
|
||||
}) {
|
||||
return (
|
||||
<div className='space-y-1.5'>
|
||||
<p className='text-xs font-medium text-muted-foreground'>{label}</p>
|
||||
{children}
|
||||
</div>
|
||||
);
|
||||
}
|
||||
|
||||
function TimeSelect({
|
||||
value,
|
||||
options,
|
||||
onValueChange,
|
||||
}: {
|
||||
value: string;
|
||||
options: string[];
|
||||
onValueChange: (value: string) => void;
|
||||
}) {
|
||||
return (
|
||||
<Select value={value} onValueChange={onValueChange}>
|
||||
<SelectTrigger className='h-8 w-18 text-xs'>
|
||||
<SelectValue />
|
||||
</SelectTrigger>
|
||||
<SelectContent>
|
||||
{options.map((option) => (
|
||||
<SelectItem key={option} value={option}>
|
||||
{option}
|
||||
</SelectItem>
|
||||
))}
|
||||
</SelectContent>
|
||||
</Select>
|
||||
);
|
||||
}
|
||||
|
||||
/** shadcn 日期 + 时间选择器,value 为 ISO 字符串。 */
|
||||
function DateTimePicker({
|
||||
value,
|
||||
onChange,
|
||||
}: {
|
||||
value: string;
|
||||
onChange: (value: string) => void;
|
||||
}) {
|
||||
const [open, setOpen] = useState(false);
|
||||
const current = value ? new Date(value) : undefined;
|
||||
|
||||
const applyDate = (date: Date | undefined) => {
|
||||
if (!date) return;
|
||||
const next = value ? new Date(value) : new Date();
|
||||
next.setFullYear(date.getFullYear(), date.getMonth(), date.getDate());
|
||||
onChange(next.toISOString());
|
||||
};
|
||||
|
||||
const applyTime = (hh: string, mm: string) => {
|
||||
const next = value ? new Date(value) : new Date();
|
||||
next.setHours(Number(hh), Number(mm), 0, 0);
|
||||
onChange(next.toISOString());
|
||||
};
|
||||
|
||||
const hour = current ? String(current.getHours()).padStart(2, '0') : '00';
|
||||
const minute = current ? String(current.getMinutes()).padStart(2, '0') : '00';
|
||||
|
||||
return (
|
||||
<Popover open={open} onOpenChange={setOpen}>
|
||||
<PopoverTrigger asChild>
|
||||
<Button
|
||||
variant='outline'
|
||||
className='h-9 w-full justify-start gap-2 px-3 text-xs font-normal'
|
||||
>
|
||||
<CalendarIcon className='size-3.5 text-muted-foreground' />
|
||||
{current ? (
|
||||
format(current, 'yyyy-MM-dd HH:mm')
|
||||
) : (
|
||||
<span className='text-muted-foreground'>选择时间</span>
|
||||
)}
|
||||
</Button>
|
||||
</PopoverTrigger>
|
||||
<PopoverContent className='w-auto p-0' align='start'>
|
||||
<Calendar mode='single' selected={current} onSelect={applyDate} />
|
||||
<div className='flex items-center gap-1.5 border-t p-2'>
|
||||
<TimeSelect
|
||||
value={hour}
|
||||
options={HOUR_OPTIONS}
|
||||
onValueChange={(h) => applyTime(h, minute)}
|
||||
/>
|
||||
<span className='text-xs text-muted-foreground'>:</span>
|
||||
<TimeSelect
|
||||
value={minute}
|
||||
options={MINUTE_OPTIONS}
|
||||
onValueChange={(m) => applyTime(hour, m)}
|
||||
/>
|
||||
</div>
|
||||
</PopoverContent>
|
||||
</Popover>
|
||||
);
|
||||
}
|
||||
|
||||
export function AccessLogFilters({
|
||||
draft,
|
||||
pageSize,
|
||||
@@ -31,28 +153,12 @@ export function AccessLogFilters({
|
||||
onSearch,
|
||||
onReset,
|
||||
}: AccessLogFiltersProps) {
|
||||
const [moreOpen, setMoreOpen] = useState(false);
|
||||
|
||||
return (
|
||||
<div className='space-y-3'>
|
||||
<div className='grid gap-3 md:grid-cols-2 xl:grid-cols-4'>
|
||||
<div className='space-y-1.5'>
|
||||
<p className='text-xs font-medium text-muted-foreground'>节点 ID</p>
|
||||
<div className='relative'>
|
||||
<Search className='absolute left-2.5 top-2.5 size-3.5 text-muted-foreground' />
|
||||
<Input
|
||||
value={draft.nodeId}
|
||||
onChange={(e) =>
|
||||
onDraftChange({ ...draft, nodeId: e.target.value })
|
||||
}
|
||||
onKeyDown={(e) => {
|
||||
if (e.key === 'Enter') onSearch();
|
||||
}}
|
||||
placeholder='按 node_id 搜索'
|
||||
className='pl-8 h-9 text-xs'
|
||||
/>
|
||||
</div>
|
||||
</div>
|
||||
<div className='space-y-1.5'>
|
||||
<p className='text-xs font-medium text-muted-foreground'>来源 IP</p>
|
||||
<FilterField label='来源 IP'>
|
||||
<Input
|
||||
value={draft.remoteAddr}
|
||||
onChange={(e) =>
|
||||
@@ -64,9 +170,8 @@ export function AccessLogFilters({
|
||||
placeholder='按 IP 搜索'
|
||||
className='h-9 text-xs'
|
||||
/>
|
||||
</div>
|
||||
<div className='space-y-1.5'>
|
||||
<p className='text-xs font-medium text-muted-foreground'>访问域名</p>
|
||||
</FilterField>
|
||||
<FilterField label='访问域名'>
|
||||
<Input
|
||||
value={draft.host}
|
||||
onChange={(e) => onDraftChange({ ...draft, host: e.target.value })}
|
||||
@@ -76,21 +181,91 @@ export function AccessLogFilters({
|
||||
placeholder='按域名搜索'
|
||||
className='h-9 text-xs'
|
||||
/>
|
||||
</div>
|
||||
<div className='space-y-1.5'>
|
||||
<p className='text-xs font-medium text-muted-foreground'>请求路径</p>
|
||||
</FilterField>
|
||||
<FilterField label='状态码'>
|
||||
<Input
|
||||
value={draft.path}
|
||||
onChange={(e) => onDraftChange({ ...draft, path: e.target.value })}
|
||||
value={draft.statusCode}
|
||||
onChange={(e) =>
|
||||
onDraftChange({
|
||||
...draft,
|
||||
statusCode: e.target.value.replace(/\D/g, '').slice(0, 3),
|
||||
})
|
||||
}
|
||||
onKeyDown={(e) => {
|
||||
if (e.key === 'Enter') onSearch();
|
||||
}}
|
||||
placeholder='按路径搜索'
|
||||
placeholder='输入状态码,如 404'
|
||||
className='h-9 text-xs'
|
||||
/>
|
||||
</div>
|
||||
</FilterField>
|
||||
</div>
|
||||
|
||||
<Collapsible open={moreOpen} onOpenChange={setMoreOpen}>
|
||||
<CollapsibleTrigger asChild>
|
||||
<Button
|
||||
variant='ghost'
|
||||
size='sm'
|
||||
className='-ml-1 h-8 gap-1 px-1 text-xs text-muted-foreground'
|
||||
>
|
||||
更多筛选
|
||||
<ChevronDown
|
||||
className={`size-3.5 transition-transform ${
|
||||
moreOpen ? 'rotate-180' : ''
|
||||
}`}
|
||||
/>
|
||||
</Button>
|
||||
</CollapsibleTrigger>
|
||||
<CollapsibleContent>
|
||||
<div className='grid gap-3 pt-3 md:grid-cols-2 xl:grid-cols-3'>
|
||||
<FilterField label='节点 ID'>
|
||||
<div className='relative'>
|
||||
<Search className='absolute left-2.5 top-2.5 size-3.5 text-muted-foreground' />
|
||||
<Input
|
||||
value={draft.nodeId}
|
||||
onChange={(e) =>
|
||||
onDraftChange({ ...draft, nodeId: e.target.value })
|
||||
}
|
||||
onKeyDown={(e) => {
|
||||
if (e.key === 'Enter') onSearch();
|
||||
}}
|
||||
placeholder='按 node_id 搜索'
|
||||
className='pl-8 h-9 text-xs'
|
||||
/>
|
||||
</div>
|
||||
</FilterField>
|
||||
<FilterField label='请求路径'>
|
||||
<Input
|
||||
value={draft.path}
|
||||
onChange={(e) =>
|
||||
onDraftChange({ ...draft, path: e.target.value })
|
||||
}
|
||||
onKeyDown={(e) => {
|
||||
if (e.key === 'Enter') onSearch();
|
||||
}}
|
||||
placeholder='按路径搜索'
|
||||
className='h-9 text-xs'
|
||||
/>
|
||||
</FilterField>
|
||||
<FilterField label='时间范围'>
|
||||
<div className='grid grid-cols-2 gap-2'>
|
||||
<DateTimePicker
|
||||
value={draft.since}
|
||||
onChange={(value) =>
|
||||
onDraftChange({ ...draft, since: value })
|
||||
}
|
||||
/>
|
||||
<DateTimePicker
|
||||
value={draft.until}
|
||||
onChange={(value) =>
|
||||
onDraftChange({ ...draft, until: value })
|
||||
}
|
||||
/>
|
||||
</div>
|
||||
</FilterField>
|
||||
</div>
|
||||
</CollapsibleContent>
|
||||
</Collapsible>
|
||||
|
||||
<div className='flex flex-col gap-3 sm:flex-row sm:items-end sm:justify-between'>
|
||||
<div className='space-y-1.5 w-full sm:max-w-[180px]'>
|
||||
<p className='text-xs font-medium text-muted-foreground'>每页条数</p>
|
||||
|
||||
@@ -5,6 +5,9 @@ export type SearchDraft = {
|
||||
remoteAddr: string;
|
||||
host: string;
|
||||
path: string;
|
||||
statusCode: string;
|
||||
since: string;
|
||||
until: string;
|
||||
};
|
||||
|
||||
export type OverviewRangeHours = 24 | 168 | 360 | 720;
|
||||
|
||||
@@ -33,6 +33,9 @@ const emptyDraft: SearchDraft = {
|
||||
remoteAddr: '',
|
||||
host: '',
|
||||
path: '',
|
||||
statusCode: '',
|
||||
since: '',
|
||||
until: '',
|
||||
};
|
||||
|
||||
function resolveTab(value: string | null): AccessLogTab {
|
||||
@@ -111,6 +114,11 @@ function AccessLogsPageContent() {
|
||||
remote_addr: filters.remoteAddr || undefined,
|
||||
host: filters.host || undefined,
|
||||
path: filters.path || undefined,
|
||||
status_code: filters.statusCode
|
||||
? Number.parseInt(filters.statusCode, 10)
|
||||
: undefined,
|
||||
since: filters.since || undefined,
|
||||
until: filters.until || undefined,
|
||||
p: page,
|
||||
page_size: pageSize,
|
||||
sort_by: detailSortState.sortBy,
|
||||
@@ -161,6 +169,9 @@ function AccessLogsPageContent() {
|
||||
remoteAddr: draft.remoteAddr.trim(),
|
||||
host: draft.host.trim(),
|
||||
path: draft.path.trim(),
|
||||
statusCode: draft.statusCode.trim(),
|
||||
since: draft.since.trim(),
|
||||
until: draft.until.trim(),
|
||||
});
|
||||
setPage(0);
|
||||
}, [draft]);
|
||||
|
||||
@@ -1,60 +0,0 @@
|
||||
'use client';
|
||||
|
||||
import { Globe2 } from 'lucide-react';
|
||||
|
||||
import {
|
||||
Card,
|
||||
CardContent,
|
||||
CardDescription,
|
||||
CardHeader,
|
||||
CardTitle,
|
||||
} from '@/components/ui/card';
|
||||
import { Progress } from '@/components/ui/progress';
|
||||
import type { DistributionItem } from '@/lib/services/openflare';
|
||||
|
||||
import { formatCompactNumber } from './dashboard-utils';
|
||||
|
||||
export function GeoDistributionList({ items }: { items: DistributionItem[] }) {
|
||||
const sortedItems = [...items]
|
||||
.sort((left, right) => right.value - left.value)
|
||||
.slice(0, 8);
|
||||
const maxValue = sortedItems[0]?.value ?? 0;
|
||||
|
||||
return (
|
||||
<Card className='border-dashed shadow-none'>
|
||||
<CardHeader>
|
||||
<CardTitle className='text-sm font-semibold flex items-center gap-1.5'>
|
||||
<Globe2 className='size-4 text-primary' />
|
||||
来源国家分布
|
||||
</CardTitle>
|
||||
<CardDescription className='text-xs'>
|
||||
聚合最近 24 小时主要来源国家。
|
||||
</CardDescription>
|
||||
</CardHeader>
|
||||
<CardContent>
|
||||
{sortedItems.length === 0 ? (
|
||||
<div className='flex min-h-[180px] items-center justify-center text-xs text-muted-foreground'>
|
||||
暂无来源分布数据
|
||||
</div>
|
||||
) : (
|
||||
<div className='space-y-3'>
|
||||
{sortedItems.map((item) => {
|
||||
const ratio = maxValue > 0 ? (item.value / maxValue) * 100 : 0;
|
||||
return (
|
||||
<div key={item.key} className='space-y-1.5'>
|
||||
<div className='flex items-center justify-between text-xs'>
|
||||
<span className='font-medium'>{item.key || '未知'}</span>
|
||||
<span className='font-mono tabular-nums text-muted-foreground'>
|
||||
{formatCompactNumber(item.value)}
|
||||
</span>
|
||||
</div>
|
||||
<Progress value={ratio} className='h-1.5' />
|
||||
</div>
|
||||
);
|
||||
})}
|
||||
</div>
|
||||
)}
|
||||
</CardContent>
|
||||
</Card>
|
||||
);
|
||||
}
|
||||
+21
-35
@@ -3,39 +3,25 @@
|
||||
import { TrendChart } from '@/components/data/trend-chart';
|
||||
import { Card, CardContent, CardHeader, CardTitle } from '@/components/ui/card';
|
||||
import type {
|
||||
DiskIOTrendPoint,
|
||||
CapacityTrendPoint,
|
||||
NetworkTrendPoint,
|
||||
} from '@/lib/services/openflare';
|
||||
|
||||
import {
|
||||
formatBytes,
|
||||
formatBytesPerSecond,
|
||||
formatTrendHour,
|
||||
} from './dashboard-utils';
|
||||
import { formatBytes, formatPercent, formatTrendHour } from './dashboard-utils';
|
||||
|
||||
/** Backend disk points are per-hour totals; chart displays bytes/s within each hour. */
|
||||
const DISK_BUCKET_SECONDS = 3600;
|
||||
|
||||
function diskBytesToRate(bytes: number) {
|
||||
return bytes > 0 ? bytes / DISK_BUCKET_SECONDS : 0;
|
||||
}
|
||||
|
||||
function formatDiskRate(bytesPerSecond: number) {
|
||||
return formatBytesPerSecond(bytesPerSecond, 1, { zeroText: '0 B' });
|
||||
}
|
||||
|
||||
export function NetworkDiskTrendChart({
|
||||
/** 业务流量(来自访问日志)与容量趋势(节点 Agent 宿主机指标)合并展示。 */
|
||||
export function TrafficCapacityTrendChart({
|
||||
networkPoints,
|
||||
diskPoints,
|
||||
capacityPoints,
|
||||
}: {
|
||||
networkPoints: NetworkTrendPoint[];
|
||||
diskPoints: DiskIOTrendPoint[];
|
||||
capacityPoints: CapacityTrendPoint[];
|
||||
}) {
|
||||
return (
|
||||
<Card className='border-dashed shadow-none'>
|
||||
<CardHeader>
|
||||
<CardTitle className='text-sm font-semibold'>
|
||||
24 小时业务流量与宿主机磁盘
|
||||
24 小时业务流量与容量趋势
|
||||
</CardTitle>
|
||||
</CardHeader>
|
||||
<CardContent className='space-y-6'>
|
||||
@@ -66,31 +52,31 @@ export function NetworkDiskTrendChart({
|
||||
/>
|
||||
|
||||
<TrendChart
|
||||
labels={diskPoints.map((point) =>
|
||||
labels={capacityPoints.map((point) =>
|
||||
formatTrendHour(point.bucket_started_at),
|
||||
)}
|
||||
height={180}
|
||||
summaryScope='average'
|
||||
summaryHint='近 24 小时 · 宿主机磁盘 · 平均速率'
|
||||
yAxisValueFormatter={formatDiskRate}
|
||||
summaryHint='近 24 小时 · 宿主机容量 · 平均值'
|
||||
yAxisValueFormatter={formatPercent}
|
||||
series={[
|
||||
{
|
||||
label: '磁盘读',
|
||||
color: '#a78bfa',
|
||||
fillColor: 'rgba(167, 139, 250, 0.14)',
|
||||
label: '平均 CPU',
|
||||
color: '#0f766e',
|
||||
fillColor: 'rgba(15, 118, 110, 0.15)',
|
||||
variant: 'area',
|
||||
values: diskPoints.map((point) =>
|
||||
diskBytesToRate(point.disk_read_bytes),
|
||||
values: capacityPoints.map(
|
||||
(point) => point.average_cpu_usage_percent,
|
||||
),
|
||||
valueFormatter: formatDiskRate,
|
||||
valueFormatter: formatPercent,
|
||||
},
|
||||
{
|
||||
label: '磁盘写',
|
||||
color: '#fb7185',
|
||||
values: diskPoints.map((point) =>
|
||||
diskBytesToRate(point.disk_write_bytes),
|
||||
label: '平均内存',
|
||||
color: '#2563eb',
|
||||
values: capacityPoints.map(
|
||||
(point) => point.average_memory_usage_percent,
|
||||
),
|
||||
valueFormatter: formatDiskRate,
|
||||
valueFormatter: formatPercent,
|
||||
},
|
||||
]}
|
||||
/>
|
||||
@@ -15,7 +15,7 @@ import { formatTrendHour } from './dashboard-utils';
|
||||
export function TrafficTrendChart({
|
||||
points,
|
||||
title = '24 小时请求趋势',
|
||||
description = '观察整体请求量和错误量是否出现异常抬升。',
|
||||
description = '按小时拆分请求总量与 2xx/4xx/5xx 状态码请求量,判断各状态是否异常抬升。',
|
||||
}: {
|
||||
points: TrafficTrendPoint[];
|
||||
title?: string;
|
||||
@@ -43,9 +43,19 @@ export function TrafficTrendChart({
|
||||
values: points.map((point) => point.request_count),
|
||||
},
|
||||
{
|
||||
label: '错误量',
|
||||
label: '2xx 请求',
|
||||
color: '#22c55e',
|
||||
values: points.map((point) => point.status_2xx_count),
|
||||
},
|
||||
{
|
||||
label: '4xx 请求',
|
||||
color: '#f97316',
|
||||
values: points.map((point) => point.status_4xx_count),
|
||||
},
|
||||
{
|
||||
label: '5xx 请求',
|
||||
color: '#ef4444',
|
||||
values: points.map((point) => point.error_count),
|
||||
values: points.map((point) => point.status_5xx_count),
|
||||
},
|
||||
]}
|
||||
/>
|
||||
|
||||
@@ -10,14 +10,13 @@ import { Button } from '@/components/ui/button';
|
||||
import { DashboardService } from '@/lib/services/openflare';
|
||||
import { formatDateTime } from '@/lib/utils';
|
||||
|
||||
import { CapacityTrendChart } from './components/dashboard/capacity-trend-chart';
|
||||
import {
|
||||
SourceDistributionChart,
|
||||
StatusCodeDistributionChart,
|
||||
TopDomainChart,
|
||||
} from './components/dashboard/distribution-rank-charts';
|
||||
import { NetworkDiskTrendChart } from './components/dashboard/network-disk-trend-chart';
|
||||
import { NodeHealthTable } from './components/dashboard/node-health-table';
|
||||
import { TrafficCapacityTrendChart } from './components/dashboard/traffic-capacity-trend-chart';
|
||||
import { TrafficTrendChart } from './components/dashboard/traffic-trend-chart';
|
||||
import { WorldStage } from './components/dashboard/world-stage';
|
||||
import { getErrorMessage } from './nodes/components/node-utils';
|
||||
@@ -84,16 +83,8 @@ export default function OpenFlareDashboardPage() {
|
||||
sourceCountries={overview.distributions.source_countries}
|
||||
/>
|
||||
|
||||
<div className='grid gap-6 xl:grid-cols-2'>
|
||||
<TrafficTrendChart points={overview.trends.traffic_24h} />
|
||||
<CapacityTrendChart points={overview.trends.capacity_24h} />
|
||||
</div>
|
||||
|
||||
<div className='grid gap-6'>
|
||||
<NetworkDiskTrendChart
|
||||
networkPoints={overview.trends.network_24h}
|
||||
diskPoints={overview.trends.disk_io_24h}
|
||||
/>
|
||||
<TrafficTrendChart points={overview.trends.traffic_24h} />
|
||||
</div>
|
||||
|
||||
<div className='grid gap-6 xl:grid-cols-3'>
|
||||
@@ -106,6 +97,13 @@ export default function OpenFlareDashboardPage() {
|
||||
<TopDomainChart items={overview.distributions.top_domains} />
|
||||
</div>
|
||||
|
||||
<div className='grid gap-6'>
|
||||
<TrafficCapacityTrendChart
|
||||
networkPoints={overview.trends.network_24h}
|
||||
capacityPoints={overview.trends.capacity_24h}
|
||||
/>
|
||||
</div>
|
||||
|
||||
<NodeHealthTable nodes={overview.nodes} />
|
||||
</>
|
||||
)}
|
||||
|
||||
@@ -55,7 +55,6 @@ import {
|
||||
LogOut,
|
||||
Settings,
|
||||
ShieldCheck,
|
||||
Terminal,
|
||||
UserRound,
|
||||
} from 'lucide-react';
|
||||
|
||||
|
||||
@@ -89,6 +89,9 @@ function normalizeTrafficTrendPoints(
|
||||
request_count: Number(item[1] ?? 0),
|
||||
error_count: Number(item[2] ?? 0),
|
||||
unique_visitor_count: Number(item[3] ?? 0),
|
||||
status_2xx_count: Number(item[4] ?? 0),
|
||||
status_4xx_count: Number(item[5] ?? 0),
|
||||
status_5xx_count: Number(item[6] ?? 0),
|
||||
}
|
||||
: item,
|
||||
);
|
||||
|
||||
@@ -630,6 +630,9 @@ export interface AccessLogFilters {
|
||||
remote_addr?: string;
|
||||
host?: string;
|
||||
path?: string;
|
||||
status_code?: number;
|
||||
since?: string;
|
||||
until?: string;
|
||||
p?: number;
|
||||
page_size?: number;
|
||||
sort_by?: string;
|
||||
@@ -1156,6 +1159,9 @@ export interface TrafficTrendPoint {
|
||||
request_count: number;
|
||||
error_count: number;
|
||||
unique_visitor_count: number;
|
||||
status_2xx_count: number;
|
||||
status_4xx_count: number;
|
||||
status_5xx_count: number;
|
||||
}
|
||||
|
||||
export interface CapacityTrendPoint {
|
||||
@@ -1224,7 +1230,15 @@ export interface DashboardOverview {
|
||||
nodes: DashboardNodeHealth[];
|
||||
}
|
||||
|
||||
export type CompactTrafficTrendPoint = [string, number, number, number];
|
||||
export type CompactTrafficTrendPoint = [
|
||||
string,
|
||||
number,
|
||||
number,
|
||||
number,
|
||||
number,
|
||||
number,
|
||||
number,
|
||||
];
|
||||
export type CompactCapacityTrendPoint = [string, number, number, number];
|
||||
/** [bucket, bytes_received, bytes_provided, reported_nodes] */
|
||||
export type CompactNetworkTrendPoint = [string, number, number, number];
|
||||
|
||||
@@ -4559,4 +4559,4 @@
|
||||
}
|
||||
}
|
||||
}
|
||||
]
|
||||
]
|
||||
|
||||
@@ -330,11 +330,15 @@ func compressDistributionItems(items []observability.DistributionItem) [][]any {
|
||||
func compressTrafficTrendPoints(points []observability.TrafficTrendPoint) [][]any {
|
||||
rows := make([][]any, 0, len(points))
|
||||
for _, point := range points {
|
||||
// Compact layout: [0] bucket, [1] request, [2] error, [3] uv, [4] 2xx, [5] 4xx, [6] 5xx
|
||||
rows = append(rows, []any{
|
||||
point.BucketStartedAt,
|
||||
point.RequestCount,
|
||||
point.ErrorCount,
|
||||
point.UniqueVisitorCount,
|
||||
point.Status2xxCount,
|
||||
point.Status4xxCount,
|
||||
point.Status5xxCount,
|
||||
})
|
||||
}
|
||||
return rows
|
||||
|
||||
@@ -143,7 +143,7 @@ func TestGetOverviewStructure(t *testing.T) {
|
||||
require.Len(t, overview.Trends.Network24h, 24)
|
||||
require.Len(t, overview.Trends.DiskIO24h, 24)
|
||||
for _, row := range overview.Trends.Traffic24h {
|
||||
require.Len(t, row, 4)
|
||||
require.Len(t, row, 7)
|
||||
}
|
||||
for _, row := range overview.Trends.Capacity24h {
|
||||
require.Len(t, row, 4)
|
||||
|
||||
@@ -39,6 +39,9 @@ type AccessLogQuery struct {
|
||||
RemoteAddr string `json:"remote_addr"`
|
||||
Host string `json:"host"`
|
||||
Path string `json:"path"`
|
||||
StatusCode int `json:"status_code"`
|
||||
Since string `json:"since"`
|
||||
Until string `json:"until"`
|
||||
Page int `json:"page"`
|
||||
PageSize int `json:"page_size"`
|
||||
SortBy string `json:"sort_by"`
|
||||
@@ -533,7 +536,10 @@ func buildAccessLogOverviewTrends(
|
||||
// ListAccessLogs returns paginated access logs.
|
||||
func ListAccessLogs(ctx context.Context, input AccessLogQuery) (*AccessLogList, error) {
|
||||
normalized := normalizeAccessLogQuery(input)
|
||||
modelQuery := buildModelAccessLogQuery(normalized)
|
||||
modelQuery, err := buildModelAccessLogQuery(normalized)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
logs, err := repository.ListOpenFlareAccessLogs(ctx, modelQuery)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -586,13 +592,17 @@ func ListFoldedAccessLogs(ctx context.Context, input AccessLogQuery) (*FoldedAcc
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
modelQuery := buildModelAccessLogQuery(normalized)
|
||||
modelQuery, err := buildModelAccessLogQuery(normalized)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
bucketQuery := model.OpenFlareAccessLogBucketQuery{
|
||||
NodeID: modelQuery.NodeID,
|
||||
RemoteAddr: modelQuery.RemoteAddr,
|
||||
Host: modelQuery.Host,
|
||||
Path: modelQuery.Path,
|
||||
Since: modelQuery.Since,
|
||||
Until: modelQuery.Until,
|
||||
Page: normalized.Page,
|
||||
PageSize: normalized.PageSize,
|
||||
SortBy: normalizeFoldSortBy(input.SortBy),
|
||||
@@ -881,18 +891,49 @@ func CleanupAccessLogs(ctx context.Context, input AccessLogCleanupInput) (*Acces
|
||||
}, nil
|
||||
}
|
||||
|
||||
func buildModelAccessLogQuery(input AccessLogQuery) model.OpenFlareAccessLogQuery {
|
||||
func buildModelAccessLogQuery(input AccessLogQuery) (model.OpenFlareAccessLogQuery, error) {
|
||||
since := defaultAccessLogSince()
|
||||
until := time.Now().UTC()
|
||||
if input.Since != "" || input.Until != "" {
|
||||
parsedSince, parsedUntil, err := resolveAccessLogWindow(input.Since, input.Until)
|
||||
if err != nil {
|
||||
return model.OpenFlareAccessLogQuery{}, err
|
||||
}
|
||||
since, until = parsedSince, parsedUntil
|
||||
}
|
||||
return model.OpenFlareAccessLogQuery{
|
||||
NodeID: strings.TrimSpace(input.NodeID),
|
||||
RemoteAddr: strings.TrimSpace(input.RemoteAddr),
|
||||
Host: strings.TrimSpace(input.Host),
|
||||
Path: strings.TrimSpace(input.Path),
|
||||
Since: defaultAccessLogSince(),
|
||||
StatusCode: input.StatusCode,
|
||||
Since: since,
|
||||
Until: until,
|
||||
Page: input.Page,
|
||||
PageSize: input.PageSize,
|
||||
SortBy: input.SortBy,
|
||||
SortOrder: input.SortOrder,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// resolveAccessLogWindow 校验并解析 RFC3339 时间窗;since/until 必须成对提供,
|
||||
// 且 until 需晚于 since。
|
||||
func resolveAccessLogWindow(sinceRaw, untilRaw string) (time.Time, time.Time, error) {
|
||||
if sinceRaw == "" || untilRaw == "" {
|
||||
return time.Time{}, time.Time{}, errors.New("since 与 until 需同时提供")
|
||||
}
|
||||
parsedSince, err := time.Parse(time.RFC3339, sinceRaw)
|
||||
if err != nil {
|
||||
return time.Time{}, time.Time{}, errors.New("since 时间格式无效,需为 RFC3339")
|
||||
}
|
||||
parsedUntil, err := time.Parse(time.RFC3339, untilRaw)
|
||||
if err != nil {
|
||||
return time.Time{}, time.Time{}, errors.New("until 时间格式无效,需为 RFC3339")
|
||||
}
|
||||
if !parsedUntil.After(parsedSince) {
|
||||
return time.Time{}, time.Time{}, errors.New("until 必须晚于 since")
|
||||
}
|
||||
return parsedSince, parsedUntil, nil
|
||||
}
|
||||
|
||||
func defaultAccessLogSince() time.Time {
|
||||
@@ -932,6 +973,9 @@ func normalizeAccessLogQuery(input AccessLogQuery) AccessLogQuery {
|
||||
RemoteAddr: strings.TrimSpace(input.RemoteAddr),
|
||||
Host: strings.TrimSpace(input.Host),
|
||||
Path: strings.TrimSpace(input.Path),
|
||||
StatusCode: input.StatusCode,
|
||||
Since: strings.TrimSpace(input.Since),
|
||||
Until: strings.TrimSpace(input.Until),
|
||||
Page: normalizeAccessLogPage(input.Page),
|
||||
PageSize: normalizeAccessLogPageSize(input.PageSize),
|
||||
SortBy: normalizeAccessLogSortBy(input.SortBy),
|
||||
|
||||
@@ -85,6 +85,9 @@ type TrafficTrendPoint struct {
|
||||
RequestCount int64 `json:"request_count"`
|
||||
ErrorCount int64 `json:"error_count"`
|
||||
UniqueVisitorCount int64 `json:"unique_visitor_count"`
|
||||
Status2xxCount int64 `json:"status_2xx_count"`
|
||||
Status4xxCount int64 `json:"status_4xx_count"`
|
||||
Status5xxCount int64 `json:"status_5xx_count"`
|
||||
}
|
||||
|
||||
// CapacityTrendPoint is a capacity trend bucket.
|
||||
@@ -303,9 +306,10 @@ func BuildNodeTrends(
|
||||
}
|
||||
}
|
||||
|
||||
// BuildTrafficTrendPointsFromAccessLogs builds 24h request/error buckets from access logs.
|
||||
// Prefers of_access_log_hourly when available; falls back to raw bucket aggregates.
|
||||
// UniqueVisitorCount on hourly path is 0 (use TrafficSummary for exact UV).
|
||||
// BuildTrafficTrendPointsFromAccessLogs builds 24h request/error/status buckets from access logs.
|
||||
// Uses raw bucket aggregates: the hourly rollup (of_access_log_hourly) has no per-status counts,
|
||||
// and the 24h window on the dashboard is cached, so the raw scan is acceptable.
|
||||
// UniqueVisitorCount from buckets is exact (uniqExact on raw); TrafficSummary is used elsewhere for UV.
|
||||
func BuildTrafficTrendPointsFromAccessLogs(ctx context.Context, now time.Time, nodeID string, since time.Time) []TrafficTrendPoint {
|
||||
start := trendWindowStart(now)
|
||||
points := make([]TrafficTrendPoint, observabilityTrendBuckets)
|
||||
@@ -313,22 +317,6 @@ func BuildTrafficTrendPointsFromAccessLogs(ctx context.Context, now time.Time, n
|
||||
points[index].BucketStartedAt = start.Add(time.Duration(index) * time.Hour)
|
||||
}
|
||||
|
||||
if hourly, err := repository.ListOpenFlareTrafficHourlySince(ctx, nodeID, since); err == nil && len(hourly) > 0 {
|
||||
for _, row := range hourly {
|
||||
if row == nil {
|
||||
continue
|
||||
}
|
||||
index, ok := trendBucketIndex(row.Hour, start)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
points[index].RequestCount += row.RequestCount
|
||||
points[index].ErrorCount += row.ErrorCount
|
||||
// UniqueVisitorCount intentionally not summed from hourly rollup (always 0 / overcounts).
|
||||
}
|
||||
return points
|
||||
}
|
||||
|
||||
buckets, err := repository.ListOpenFlareAccessLogBuckets(ctx, model.OpenFlareAccessLogBucketQuery{
|
||||
NodeID: nodeID,
|
||||
Since: since,
|
||||
@@ -353,6 +341,9 @@ func BuildTrafficTrendPointsFromAccessLogs(ctx context.Context, now time.Time, n
|
||||
points[index].RequestCount = row.RequestCount
|
||||
points[index].ErrorCount = row.ServerErrorCount
|
||||
points[index].UniqueVisitorCount = row.UniqueIPCount
|
||||
points[index].Status2xxCount = row.Status2xxCount
|
||||
points[index].Status4xxCount = row.Status4xxCount
|
||||
points[index].Status5xxCount = row.Status5xxCount
|
||||
}
|
||||
}
|
||||
return points
|
||||
|
||||
@@ -0,0 +1,9 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package observability defines shared error messages for observability operations.
|
||||
package observability
|
||||
|
||||
const (
|
||||
errInvalidStatusCode = "status_code 必须为 100-599 之间的整数"
|
||||
)
|
||||
@@ -55,3 +55,76 @@ func TestReadQueryStringArrayAcceptsHostsBracketForm(t *testing.T) {
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadAccessLogQueryIncludesStatusCode(t *testing.T) {
|
||||
gin.SetMode(gin.TestMode)
|
||||
w := httptest.NewRecorder()
|
||||
c, _ := gin.CreateTestContext(w)
|
||||
req, err := http.NewRequest(
|
||||
http.MethodGet,
|
||||
"/?node_id=n1&remote_addr=1.2.3.4&host=a.example&path=/api&status_code=404&p=2&page_size=50",
|
||||
nil,
|
||||
)
|
||||
require.NoError(t, err)
|
||||
c.Request = req
|
||||
|
||||
got, err := readAccessLogQuery(c)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, "n1", got.NodeID)
|
||||
require.Equal(t, "1.2.3.4", got.RemoteAddr)
|
||||
require.Equal(t, "a.example", got.Host)
|
||||
require.Equal(t, "/api", got.Path)
|
||||
require.Equal(t, 404, got.StatusCode)
|
||||
require.Equal(t, 2, got.Page)
|
||||
require.Equal(t, 50, got.PageSize)
|
||||
}
|
||||
|
||||
func TestReadAccessLogQueryRejectsInvalidStatusCode(t *testing.T) {
|
||||
gin.SetMode(gin.TestMode)
|
||||
for _, raw := range []string{"abc", "99", "600", "-1"} {
|
||||
w := httptest.NewRecorder()
|
||||
c, _ := gin.CreateTestContext(w)
|
||||
req, err := http.NewRequest(http.MethodGet, "/?status_code="+raw, nil)
|
||||
require.NoError(t, err)
|
||||
c.Request = req
|
||||
|
||||
_, err = readAccessLogQuery(c)
|
||||
require.Error(t, err, "status_code=%s should be rejected", raw)
|
||||
}
|
||||
|
||||
w := httptest.NewRecorder()
|
||||
c, _ := gin.CreateTestContext(w)
|
||||
req, err := http.NewRequest(http.MethodGet, "/", nil)
|
||||
require.NoError(t, err)
|
||||
c.Request = req
|
||||
got, err := readAccessLogQuery(c)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 0, got.StatusCode)
|
||||
}
|
||||
|
||||
func TestResolveAccessLogWindow(t *testing.T) {
|
||||
since, until, err := resolveAccessLogWindow(
|
||||
"2026-08-01T00:00:00Z",
|
||||
"2026-08-02T00:00:00Z",
|
||||
)
|
||||
require.NoError(t, err)
|
||||
require.True(t, until.After(since))
|
||||
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
since string
|
||||
until string
|
||||
}{
|
||||
{name: "missing both", since: "", until: ""},
|
||||
{name: "only since", since: "2026-08-01T00:00:00Z", until: ""},
|
||||
{name: "only until", since: "", until: "2026-08-02T00:00:00Z"},
|
||||
{name: "bad since", since: "not-a-time", until: "2026-08-02T00:00:00Z"},
|
||||
{name: "bad until", since: "2026-08-01T00:00:00Z", until: "not-a-time"},
|
||||
{name: "reversed", since: "2026-08-02T00:00:00Z", until: "2026-08-01T00:00:00Z"},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
_, _, err := resolveAccessLogWindow(tc.since, tc.until)
|
||||
require.Error(t, err)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@
|
||||
package observability
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"net/http"
|
||||
"strconv"
|
||||
|
||||
@@ -45,7 +46,7 @@ func GetAccessLogOverviewHandler(c *gin.Context) {
|
||||
|
||||
// GetAccessLogsHandler 分页列出访问日志。
|
||||
// @Summary 列出访问日志
|
||||
// @Description 分页返回 OpenFlare 访问日志,支持按节点、IP、主机与路径筛选,需要管理员权限
|
||||
// @Description 分页返回 OpenFlare 访问日志,支持按节点、IP、主机、路径与状态码筛选,需要管理员权限
|
||||
// @Tags openflare-observability
|
||||
// @Produce json
|
||||
// @Security SessionCookie
|
||||
@@ -53,6 +54,9 @@ func GetAccessLogOverviewHandler(c *gin.Context) {
|
||||
// @Param remote_addr query string false "客户端 IP"
|
||||
// @Param host query string false "请求 Host"
|
||||
// @Param path query string false "请求路径"
|
||||
// @Param status_code query int false "HTTP 状态码(100-599)"
|
||||
// @Param since query string false "起始时间(RFC3339,需与 until 成对提供)"
|
||||
// @Param until query string false "结束时间(RFC3339,需与 since 成对提供)"
|
||||
// @Param p query int false "页码"
|
||||
// @Param page_size query int false "每页条数"
|
||||
// @Param sort_by query string false "排序字段"
|
||||
@@ -64,7 +68,11 @@ func GetAccessLogOverviewHandler(c *gin.Context) {
|
||||
// @Failure 500 {object} response.Any "内部错误"
|
||||
// @Router /api/v1/d/access-logs [get]
|
||||
func GetAccessLogsHandler(c *gin.Context) {
|
||||
logs, err := ListAccessLogs(c.Request.Context(), readAccessLogQuery(c))
|
||||
query, err := readAccessLogQuery(c)
|
||||
if apiutil.AbortBadRequestOnError(c, err) {
|
||||
return
|
||||
}
|
||||
logs, err := ListAccessLogs(c.Request.Context(), query)
|
||||
if apiutil.AbortBadRequestOnError(c, err) {
|
||||
return
|
||||
}
|
||||
@@ -93,7 +101,10 @@ func GetAccessLogsHandler(c *gin.Context) {
|
||||
// @Failure 500 {object} response.Any "内部错误"
|
||||
// @Router /api/v1/d/access-logs/folds [get]
|
||||
func GetFoldedAccessLogsHandler(c *gin.Context) {
|
||||
query := readAccessLogQuery(c)
|
||||
query, err := readAccessLogQuery(c)
|
||||
if apiutil.AbortBadRequestOnError(c, err) {
|
||||
return
|
||||
}
|
||||
query.FoldMinutes = readQueryInt(c, "fold_minutes")
|
||||
logs, err := ListFoldedAccessLogs(c.Request.Context(), query)
|
||||
if apiutil.AbortBadRequestOnError(c, err) {
|
||||
@@ -270,17 +281,27 @@ func CleanupAccessLogsHandler(c *gin.Context) {
|
||||
c.JSON(http.StatusOK, response.OK(result))
|
||||
}
|
||||
|
||||
func readAccessLogQuery(c *gin.Context) AccessLogQuery {
|
||||
return AccessLogQuery{
|
||||
func readAccessLogQuery(c *gin.Context) (AccessLogQuery, error) {
|
||||
query := AccessLogQuery{
|
||||
NodeID: c.Query("node_id"),
|
||||
RemoteAddr: c.Query("remote_addr"),
|
||||
Host: c.Query("host"),
|
||||
Path: c.Query("path"),
|
||||
Since: c.Query("since"),
|
||||
Until: c.Query("until"),
|
||||
Page: readQueryInt(c, "p"),
|
||||
PageSize: readQueryInt(c, "page_size"),
|
||||
SortBy: c.Query("sort_by"),
|
||||
SortOrder: c.Query("sort_order"),
|
||||
}
|
||||
if raw := c.Query("status_code"); raw != "" {
|
||||
code, err := strconv.Atoi(raw)
|
||||
if err != nil || code < 100 || code > 599 {
|
||||
return AccessLogQuery{}, errors.New(errInvalidStatusCode)
|
||||
}
|
||||
query.StatusCode = code
|
||||
}
|
||||
return query, nil
|
||||
}
|
||||
|
||||
func readQueryInt(c *gin.Context, key string) int {
|
||||
|
||||
@@ -25,14 +25,16 @@ type NodeAccessLogFilter struct {
|
||||
RemoteAddr string
|
||||
Host string
|
||||
// Hosts exact-matches any host (case-insensitive). Prefer over Host for multi-domain scopes.
|
||||
Hosts []string
|
||||
Path string
|
||||
Since time.Time
|
||||
Until time.Time
|
||||
Page int
|
||||
PageSize int
|
||||
SortBy string
|
||||
SortOrder string
|
||||
Hosts []string
|
||||
Path string
|
||||
// StatusCode filters by exact HTTP status code when > 0.
|
||||
StatusCode int
|
||||
Since time.Time
|
||||
Until time.Time
|
||||
Page int
|
||||
PageSize int
|
||||
SortBy string
|
||||
SortOrder string
|
||||
}
|
||||
|
||||
// NodeObservabilityFilter scopes ClickHouse node observability queries.
|
||||
|
||||
@@ -10,6 +10,9 @@ type NodeAccessLogBucketAggregate struct {
|
||||
SuccessCount int64 `gorm:"column:success_count"`
|
||||
ClientErrorCount int64 `gorm:"column:client_error_count"`
|
||||
ServerErrorCount int64 `gorm:"column:server_error_count"`
|
||||
Status2xxCount int64 `gorm:"column:status_2xx_count"`
|
||||
Status4xxCount int64 `gorm:"column:status_4xx_count"`
|
||||
Status5xxCount int64 `gorm:"column:status_5xx_count"`
|
||||
UniqueIPCount int64 `gorm:"column:unique_ip_count"`
|
||||
UniqueHostCount int64 `gorm:"column:unique_host_count"`
|
||||
BytesSent int64 `gorm:"column:bytes_sent"`
|
||||
|
||||
@@ -161,14 +161,16 @@ type OpenFlareAccessLogQuery struct {
|
||||
RemoteAddr string
|
||||
Host string
|
||||
// Hosts exact-matches any host (case-insensitive). Prefer over Host for multi-domain scopes.
|
||||
Hosts []string
|
||||
Path string
|
||||
Since time.Time
|
||||
Until time.Time
|
||||
Page int
|
||||
PageSize int
|
||||
SortBy string
|
||||
SortOrder string
|
||||
Hosts []string
|
||||
Path string
|
||||
// StatusCode filters by exact HTTP status code when > 0.
|
||||
StatusCode int
|
||||
Since time.Time
|
||||
Until time.Time
|
||||
Page int
|
||||
PageSize int
|
||||
SortBy string
|
||||
SortOrder string
|
||||
}
|
||||
|
||||
// OpenFlareAccessLogBucketQuery filters folded access log queries (v1 stub).
|
||||
@@ -196,6 +198,9 @@ type OpenFlareAccessLogBucketRow struct {
|
||||
SuccessCount int64 `json:"success_count"`
|
||||
ClientErrorCount int64 `json:"client_error_count"`
|
||||
ServerErrorCount int64 `json:"server_error_count"`
|
||||
Status2xxCount int64 `json:"status_2xx_count"`
|
||||
Status4xxCount int64 `json:"status_4xx_count"`
|
||||
Status5xxCount int64 `json:"status_5xx_count"`
|
||||
BytesSent int64 `json:"bytes_sent"`
|
||||
RequestLength int64 `json:"request_length"`
|
||||
}
|
||||
|
||||
@@ -45,6 +45,16 @@ func TestBuildUserAccessLogFilterClause_EmptyUserIDs(t *testing.T) {
|
||||
assert.False(t, ok)
|
||||
}
|
||||
|
||||
func TestBuildNodeAccessLogFilterClause_StatusCode(t *testing.T) {
|
||||
clause, args := buildNodeAccessLogFilterClause(NodeAccessLogFilter{StatusCode: 404})
|
||||
assert.Equal(t, "status_code = ?", clause)
|
||||
assert.Equal(t, []any{404}, args)
|
||||
|
||||
clause, args = buildNodeAccessLogFilterClause(NodeAccessLogFilter{})
|
||||
assert.Equal(t, "1", clause)
|
||||
assert.Nil(t, args)
|
||||
}
|
||||
|
||||
func TestCountAccessLogs_EmptyUserIDs(t *testing.T) {
|
||||
count, err := CountAccessLogs(context.Background(), AccessLogFilter{UserIDs: []uint64{}})
|
||||
require.NoError(t, err)
|
||||
|
||||
@@ -11,7 +11,7 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
nodeAccessLogFilterClauseCapacity = 6
|
||||
nodeAccessLogFilterClauseCapacity = 7
|
||||
|
||||
nodeAccessLogSortDesc = "DESC"
|
||||
nodeAccessLogSortAsc = "ASC"
|
||||
@@ -55,6 +55,10 @@ func buildNodeAccessLogFilterClause(filter NodeAccessLogFilter) (string, []any)
|
||||
parts = append(parts, "path LIKE ?")
|
||||
args = append(args, trimmed+"%")
|
||||
}
|
||||
if filter.StatusCode > 0 {
|
||||
parts = append(parts, "status_code = ?")
|
||||
args = append(args, filter.StatusCode)
|
||||
}
|
||||
if !filter.Since.IsZero() {
|
||||
parts = append(parts, "logged_at >= ?")
|
||||
args = append(args, filter.Since.UTC())
|
||||
|
||||
@@ -46,6 +46,9 @@ SELECT
|
||||
countIf(status_code < 400) AS success_count,
|
||||
countIf(status_code >= 400 AND status_code < 500) AS client_error_count,
|
||||
countIf(status_code >= 500) AS server_error_count,
|
||||
countIf(status_code >= 200 AND status_code < 300) AS status_2xx_count,
|
||||
countIf(status_code >= 400 AND status_code < 500) AS status_4xx_count,
|
||||
countIf(status_code >= 500) AS status_5xx_count,
|
||||
uniqExactIf(remote_addr, remote_addr != '') AS unique_ip_count,
|
||||
uniqExactIf(host, host != '') AS unique_host_count,
|
||||
sum(bytes_sent) AS bytes_sent,
|
||||
@@ -70,10 +73,10 @@ ORDER BY %s`, bucketExpr, tableName, clause, nodeAccessLogBucketOrderClause(filt
|
||||
var result []NodeAccessLogBucketAggregate
|
||||
for rows.Next() {
|
||||
var (
|
||||
bucketEpoch int64
|
||||
requestCount, successCount, clientErrorCount, serverErrorCount, uniqueIPCount, uniqueHostCount, bytesSent, requestLength uint64
|
||||
bucketEpoch int64
|
||||
requestCount, successCount, clientErrorCount, serverErrorCount, status2xxCount, status4xxCount, status5xxCount, uniqueIPCount, uniqueHostCount, bytesSent, requestLength uint64
|
||||
)
|
||||
if err := rows.Scan(&bucketEpoch, &requestCount, &successCount, &clientErrorCount, &serverErrorCount, &uniqueIPCount, &uniqueHostCount, &bytesSent, &requestLength); err != nil {
|
||||
if err := rows.Scan(&bucketEpoch, &requestCount, &successCount, &clientErrorCount, &serverErrorCount, &status2xxCount, &status4xxCount, &status5xxCount, &uniqueIPCount, &uniqueHostCount, &bytesSent, &requestLength); err != nil {
|
||||
return nil, fmt.Errorf("scan bucket aggregate row: %w", err)
|
||||
}
|
||||
result = append(result, NodeAccessLogBucketAggregate{
|
||||
@@ -82,6 +85,9 @@ ORDER BY %s`, bucketExpr, tableName, clause, nodeAccessLogBucketOrderClause(filt
|
||||
SuccessCount: safeInt64Count(successCount),
|
||||
ClientErrorCount: safeInt64Count(clientErrorCount),
|
||||
ServerErrorCount: safeInt64Count(serverErrorCount),
|
||||
Status2xxCount: safeInt64Count(status2xxCount),
|
||||
Status4xxCount: safeInt64Count(status4xxCount),
|
||||
Status5xxCount: safeInt64Count(status5xxCount),
|
||||
UniqueIPCount: safeInt64Count(uniqueIPCount),
|
||||
UniqueHostCount: safeInt64Count(uniqueHostCount),
|
||||
BytesSent: safeInt64Count(bytesSent),
|
||||
|
||||
@@ -459,6 +459,159 @@ func TestDropExpiredPartitionsTimezoneSafety(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestBatchInsertGeneratesIDsPostgres 回归:PG 日志表 id BIGINT NOT NULL 且无默认值;
|
||||
// GORM 把零值 uint64 主键视为自增并省略 id 列,直接插入会报 23502 not-null 违例。
|
||||
// 验证 6 张日志表 BatchInsert* 为零 ID 行生成雪花 ID 后正常落库(修复前本测试失败)。
|
||||
func TestBatchInsertGeneratesIDsPostgres(t *testing.T) {
|
||||
dsn := strings.TrimSpace(os.Getenv("TEST_POSTGRES_DSN"))
|
||||
if dsn == "" {
|
||||
t.Skip("TEST_POSTGRES_DSN is not set")
|
||||
}
|
||||
|
||||
gdb, err := gorm.Open(postgres.Open(dsn), &gorm.Config{
|
||||
DisableForeignKeyConstraintWhenMigrating: true,
|
||||
Logger: logger.Default.LogMode(logger.Silent),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("open postgres: %v", err)
|
||||
}
|
||||
sqlDB, err := gdb.DB()
|
||||
if err != nil {
|
||||
t.Fatalf("sql db: %v", err)
|
||||
}
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
|
||||
schema := fmt.Sprintf("logstore_ids_%d", time.Now().UnixNano())
|
||||
if !regexp.MustCompile(`^[a-z0-9_]+$`).MatchString(schema) {
|
||||
t.Fatalf("invalid schema: %s", schema)
|
||||
}
|
||||
if err := gdb.Exec(`CREATE SCHEMA "` + schema + `"`).Error; err != nil {
|
||||
t.Fatalf("create schema: %v", err)
|
||||
}
|
||||
if err := gdb.Exec(`SET search_path TO "` + schema + `"`).Error; err != nil {
|
||||
t.Fatalf("set search_path: %v", err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
_ = gdb.Exec("SET search_path TO public").Error
|
||||
_ = gdb.Exec(`DROP SCHEMA IF EXISTS "` + schema + `" CASCADE`).Error
|
||||
_ = sqlDB.Close()
|
||||
})
|
||||
|
||||
for _, ddl := range []string{
|
||||
postgresNodeAccessLogsDDL,
|
||||
postgresUserAccessLogsDDL,
|
||||
postgresMetricSnapshotsDDL,
|
||||
postgresEdgeHealthDDL,
|
||||
postgresObsFrpsDDL,
|
||||
postgresObsFrpcDDL,
|
||||
} {
|
||||
if err := gdb.Exec(ddl).Error; err != nil {
|
||||
t.Fatalf("create table: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
ResetForTest()
|
||||
SetConfigReader(func(_ context.Context, _ string) (string, error) { return "", nil })
|
||||
defer ResetForTest()
|
||||
|
||||
ctx := context.Background()
|
||||
store := newGormStore(gdb)
|
||||
ua := newUserAccessLogGormStore(gdb)
|
||||
|
||||
now := time.Now().UTC()
|
||||
if err := store.EnsurePartitions(ctx, now, now.AddDate(0, 1, 0)); err != nil {
|
||||
t.Fatalf("EnsurePartitions: %v", err)
|
||||
}
|
||||
|
||||
nodeRows := []analyticsmodel.NodeAccessLog{
|
||||
{NodeID: "n1", LoggedAt: now, RemoteAddr: "1.1.1.1", StatusCode: 200},
|
||||
{NodeID: "n1", LoggedAt: now.Add(time.Second), RemoteAddr: "2.2.2.2", StatusCode: 500},
|
||||
}
|
||||
if err := store.BatchInsertNodeAccessLogs(ctx, nodeRows); err != nil {
|
||||
t.Fatalf("insert node access logs with zero ids: %v", err)
|
||||
}
|
||||
if nodeRows[0].ID == 0 || nodeRows[1].ID == 0 || nodeRows[0].ID == nodeRows[1].ID {
|
||||
t.Fatalf("node access log ids not generated: %+v", nodeRows)
|
||||
}
|
||||
|
||||
metricRows := []analyticsmodel.NodeMetricSnapshot{
|
||||
{NodeID: "n1", CapturedAt: now},
|
||||
{NodeID: "n2", CapturedAt: now},
|
||||
}
|
||||
if err := store.BatchInsertNodeMetricSnapshots(ctx, metricRows); err != nil {
|
||||
t.Fatalf("insert metric snapshots with zero ids: %v", err)
|
||||
}
|
||||
if metricRows[0].ID == 0 || metricRows[1].ID == 0 || metricRows[0].ID == metricRows[1].ID {
|
||||
t.Fatalf("metric snapshot ids not generated: %+v", metricRows)
|
||||
}
|
||||
|
||||
edgeRows := []analyticsmodel.NodeEdgeHealth{
|
||||
{NodeID: "n1", CapturedAt: now, Status: "ok"},
|
||||
{NodeID: "n2", CapturedAt: now, Status: "ok"},
|
||||
}
|
||||
if err := store.BatchInsertNodeEdgeHealth(ctx, edgeRows); err != nil {
|
||||
t.Fatalf("insert edge health with zero ids: %v", err)
|
||||
}
|
||||
if edgeRows[0].ID == 0 || edgeRows[1].ID == 0 || edgeRows[0].ID == edgeRows[1].ID {
|
||||
t.Fatalf("edge health ids not generated: %+v", edgeRows)
|
||||
}
|
||||
|
||||
frpsRows := []analyticsmodel.NodeObsFrps{
|
||||
{NodeID: "n1", CapturedAt: now, FrpsConnections: 1},
|
||||
{NodeID: "n2", CapturedAt: now, FrpsConnections: 2},
|
||||
}
|
||||
if err := store.BatchInsertNodeObsFrps(ctx, frpsRows); err != nil {
|
||||
t.Fatalf("insert obs frps with zero ids: %v", err)
|
||||
}
|
||||
if frpsRows[0].ID == 0 || frpsRows[1].ID == 0 || frpsRows[0].ID == frpsRows[1].ID {
|
||||
t.Fatalf("obs frps ids not generated: %+v", frpsRows)
|
||||
}
|
||||
|
||||
frpcRows := []analyticsmodel.NodeObsFrpc{
|
||||
{NodeID: "n1", CapturedAt: now, TunnelStatus: "online"},
|
||||
{NodeID: "n2", CapturedAt: now, TunnelStatus: "online"},
|
||||
}
|
||||
if err := store.BatchInsertNodeObsFrpc(ctx, frpcRows); err != nil {
|
||||
t.Fatalf("insert obs frpc with zero ids: %v", err)
|
||||
}
|
||||
if frpcRows[0].ID == 0 || frpcRows[1].ID == 0 || frpcRows[0].ID == frpcRows[1].ID {
|
||||
t.Fatalf("obs frpc ids not generated: %+v", frpcRows)
|
||||
}
|
||||
|
||||
userRows := []analyticsmodel.UserAccessLog{
|
||||
{UserID: 101, Path: "/a", CreatedAt: now},
|
||||
{UserID: 102, Path: "/b", CreatedAt: now},
|
||||
}
|
||||
if err := ua.BatchInsert(ctx, userRows); err != nil {
|
||||
t.Fatalf("insert user access logs with zero ids: %v", err)
|
||||
}
|
||||
if userRows[0].ID == 0 || userRows[1].ID == 0 || userRows[0].ID == userRows[1].ID {
|
||||
t.Fatalf("user access log ids not generated: %+v", userRows)
|
||||
}
|
||||
|
||||
expect := []struct {
|
||||
name string
|
||||
model any
|
||||
want int64
|
||||
}{
|
||||
{"of_node_access_logs", &analyticsmodel.NodeAccessLog{}, 2},
|
||||
{"of_node_metric_snapshots", &analyticsmodel.NodeMetricSnapshot{}, 2},
|
||||
{"of_node_edge_health", &analyticsmodel.NodeEdgeHealth{}, 2},
|
||||
{"of_node_obs_frps", &analyticsmodel.NodeObsFrps{}, 2},
|
||||
{"of_node_obs_frpc", &analyticsmodel.NodeObsFrpc{}, 2},
|
||||
{"w_user_access_logs", &analyticsmodel.UserAccessLog{}, 2},
|
||||
}
|
||||
for _, e := range expect {
|
||||
var got int64
|
||||
if err := gdb.Model(e.model).Count(&got).Error; err != nil {
|
||||
t.Fatalf("count %s: %v", e.name, err)
|
||||
}
|
||||
if got != e.want {
|
||||
t.Fatalf("%s count = %d, want %d", e.name, got, e.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// postgresNodeAccessLogsDDL 与 goose/postgres/202608080001_create_log_tables.sql 对齐。
|
||||
const postgresNodeAccessLogsDDL = `
|
||||
CREATE TABLE IF NOT EXISTS of_node_access_logs (
|
||||
@@ -494,3 +647,54 @@ CREATE TABLE IF NOT EXISTS w_user_access_logs (
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
||||
PRIMARY KEY (id, created_at)
|
||||
) PARTITION BY RANGE (created_at)`
|
||||
|
||||
// postgresMetricSnapshotsDDL / postgresEdgeHealthDDL / postgresObsFrpsDDL / postgresObsFrpcDDL
|
||||
// 与 goose/postgres/202608080001_create_log_tables.sql 对齐(普通表,无分区)。
|
||||
const postgresMetricSnapshotsDDL = `
|
||||
CREATE TABLE IF NOT EXISTS of_node_metric_snapshots (
|
||||
id BIGINT NOT NULL PRIMARY KEY,
|
||||
node_id VARCHAR(64) NOT NULL DEFAULT '',
|
||||
captured_at TIMESTAMPTZ NOT NULL,
|
||||
cpu_usage_percent DOUBLE PRECISION NOT NULL DEFAULT 0,
|
||||
memory_used_bytes BIGINT NOT NULL DEFAULT 0,
|
||||
memory_total_bytes BIGINT NOT NULL DEFAULT 0,
|
||||
storage_used_bytes BIGINT NOT NULL DEFAULT 0,
|
||||
storage_total_bytes BIGINT NOT NULL DEFAULT 0,
|
||||
disk_read_bytes BIGINT NOT NULL DEFAULT 0,
|
||||
disk_write_bytes BIGINT NOT NULL DEFAULT 0,
|
||||
network_rx_bytes BIGINT NOT NULL DEFAULT 0,
|
||||
network_tx_bytes BIGINT NOT NULL DEFAULT 0,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
|
||||
)`
|
||||
|
||||
const postgresEdgeHealthDDL = `
|
||||
CREATE TABLE IF NOT EXISTS of_node_edge_health (
|
||||
id BIGINT NOT NULL PRIMARY KEY,
|
||||
node_id VARCHAR(64) NOT NULL DEFAULT '',
|
||||
captured_at TIMESTAMPTZ NOT NULL,
|
||||
status VARCHAR(64) NOT NULL DEFAULT '',
|
||||
connections BIGINT NOT NULL DEFAULT 0,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
|
||||
)`
|
||||
|
||||
const postgresObsFrpsDDL = `
|
||||
CREATE TABLE IF NOT EXISTS of_node_obs_frps (
|
||||
id BIGINT NOT NULL PRIMARY KEY,
|
||||
node_id VARCHAR(64) NOT NULL DEFAULT '',
|
||||
captured_at TIMESTAMPTZ NOT NULL,
|
||||
frps_connections INTEGER NOT NULL DEFAULT 0,
|
||||
frps_proxy_count INTEGER NOT NULL DEFAULT 0,
|
||||
frps_client_count INTEGER NOT NULL DEFAULT 0,
|
||||
frps_proxies TEXT NOT NULL DEFAULT '',
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
|
||||
)`
|
||||
|
||||
const postgresObsFrpcDDL = `
|
||||
CREATE TABLE IF NOT EXISTS of_node_obs_frpc (
|
||||
id BIGINT NOT NULL PRIMARY KEY,
|
||||
node_id VARCHAR(64) NOT NULL DEFAULT '',
|
||||
captured_at TIMESTAMPTZ NOT NULL,
|
||||
tunnel_status VARCHAR(16) NOT NULL DEFAULT '',
|
||||
connected_relays_count INTEGER NOT NULL DEFAULT 0,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
|
||||
)`
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/persistence/idgen"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
|
||||
"gorm.io/gorm"
|
||||
@@ -113,6 +114,13 @@ func (s *gormLogStore) BatchInsertNodeAccessLogs(ctx context.Context, rows []ana
|
||||
if err := s.ensureWritable(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
// GORM 对零值 uint64 主键会省略 id 列;PG 日志表 id 为 NOT NULL 且无默认值,
|
||||
// 须在落库前显式生成雪花 ID(与 CH 写入路径 analytics repo BatchInsert* 行为一致)。
|
||||
for i := range rows {
|
||||
if rows[i].ID == 0 {
|
||||
rows[i].ID = idgen.NextUint64ID()
|
||||
}
|
||||
}
|
||||
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
||||
}
|
||||
|
||||
@@ -274,7 +282,12 @@ func (s *gormLogStore) RegionCounts(ctx context.Context, nodeID string, since ti
|
||||
var rows []row
|
||||
q := s.db.WithContext(ctx).Model(&analyticsmodel.NodeAccessLog{}).
|
||||
Select("region, COUNT(*) AS count").
|
||||
Where("node_id = ? AND region <> '' AND logged_at >= ?", nodeID, since)
|
||||
Where("trim(region) <> '' AND logged_at >= ?", since)
|
||||
// 空 nodeID 表示全节点聚合(对齐 CH 语义),仅非空时追加 node_id 过滤,
|
||||
// 避免 `node_id = ''` 恒空导致首页来源分布无数据。
|
||||
if nodeID = strings.TrimSpace(nodeID); nodeID != "" {
|
||||
q = q.Where("node_id = ?", nodeID)
|
||||
}
|
||||
if err := q.Group("region").Order("count DESC").Limit(limitOr(limit, defaultTopN)).Scan(&rows).Error; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -294,6 +307,9 @@ func (s *gormLogStore) BucketAggregates(ctx context.Context, query model.OpenFla
|
||||
SuccessCount int64
|
||||
ClientErrorCount int64
|
||||
ServerErrorCount int64
|
||||
Status2xxCount int64 `gorm:"column:status_2xx_count"`
|
||||
Status4xxCount int64 `gorm:"column:status_4xx_count"`
|
||||
Status5xxCount int64 `gorm:"column:status_5xx_count"`
|
||||
UniqueIPCount int64
|
||||
UniqueHostCount int64
|
||||
BytesSent int64
|
||||
@@ -309,6 +325,9 @@ func (s *gormLogStore) BucketAggregates(ctx context.Context, query model.OpenFla
|
||||
"COUNT(*) FILTER (WHERE status_code < 400) AS success_count, " +
|
||||
"COUNT(*) FILTER (WHERE status_code >= 400 AND status_code < 500) AS client_error_count, " +
|
||||
"COUNT(*) FILTER (WHERE status_code >= 500) AS server_error_count, " +
|
||||
"COUNT(*) FILTER (WHERE status_code >= 200 AND status_code < 300) AS status_2xx_count, " +
|
||||
"COUNT(*) FILTER (WHERE status_code >= 400 AND status_code < 500) AS status_4xx_count, " +
|
||||
"COUNT(*) FILTER (WHERE status_code >= 500) AS status_5xx_count, " +
|
||||
distinctNonEmptyCountSQL(s.db, "remote_addr") + " AS unique_ip_count, " +
|
||||
distinctNonEmptyCountSQL(s.db, "host") + " AS unique_host_count, " +
|
||||
"COALESCE(SUM(bytes_sent),0) AS bytes_sent, " +
|
||||
@@ -325,6 +344,9 @@ func (s *gormLogStore) BucketAggregates(ctx context.Context, query model.OpenFla
|
||||
SuccessCount: r.SuccessCount,
|
||||
ClientErrorCount: r.ClientErrorCount,
|
||||
ServerErrorCount: r.ServerErrorCount,
|
||||
Status2xxCount: r.Status2xxCount,
|
||||
Status4xxCount: r.Status4xxCount,
|
||||
Status5xxCount: r.Status5xxCount,
|
||||
UniqueIPCount: r.UniqueIPCount,
|
||||
UniqueHostCount: r.UniqueHostCount,
|
||||
BytesSent: r.BytesSent,
|
||||
@@ -780,6 +802,7 @@ func toNodeAccessLogFilter(query model.OpenFlareAccessLogQuery) analyticsmodel.N
|
||||
Host: query.Host,
|
||||
Hosts: query.Hosts,
|
||||
Path: query.Path,
|
||||
StatusCode: query.StatusCode,
|
||||
Since: query.Since,
|
||||
Until: query.Until,
|
||||
Page: query.Page,
|
||||
@@ -815,6 +838,10 @@ func buildNodeAccessLogFilterParts(f analyticsmodel.NodeAccessLogFilter) (string
|
||||
parts = append(parts, "path LIKE ?")
|
||||
args = append(args, path+"%")
|
||||
}
|
||||
if f.StatusCode > 0 {
|
||||
parts = append(parts, "status_code = ?")
|
||||
args = append(args, f.StatusCode)
|
||||
}
|
||||
if !f.Since.IsZero() {
|
||||
parts = append(parts, "logged_at >= ?")
|
||||
args = append(args, f.Since)
|
||||
@@ -1000,6 +1027,11 @@ func (s *gormLogStore) BatchInsertNodeMetricSnapshots(ctx context.Context, rows
|
||||
if err := s.ensureWritable(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
for i := range rows {
|
||||
if rows[i].ID == 0 {
|
||||
rows[i].ID = idgen.NextUint64ID()
|
||||
}
|
||||
}
|
||||
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
||||
}
|
||||
|
||||
@@ -1200,6 +1232,11 @@ func (s *gormLogStore) BatchInsertNodeEdgeHealth(ctx context.Context, rows []ana
|
||||
if err := s.ensureWritable(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
for i := range rows {
|
||||
if rows[i].ID == 0 {
|
||||
rows[i].ID = idgen.NextUint64ID()
|
||||
}
|
||||
}
|
||||
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
||||
}
|
||||
|
||||
@@ -1250,6 +1287,11 @@ func (s *gormLogStore) BatchInsertNodeObsFrps(ctx context.Context, rows []analyt
|
||||
if err := s.ensureWritable(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
for i := range rows {
|
||||
if rows[i].ID == 0 {
|
||||
rows[i].ID = idgen.NextUint64ID()
|
||||
}
|
||||
}
|
||||
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
||||
}
|
||||
|
||||
@@ -1300,6 +1342,11 @@ func (s *gormLogStore) BatchInsertNodeObsFrpc(ctx context.Context, rows []analyt
|
||||
if err := s.ensureWritable(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
for i := range rows {
|
||||
if rows[i].ID == 0 {
|
||||
rows[i].ID = idgen.NextUint64ID()
|
||||
}
|
||||
}
|
||||
return s.db.WithContext(ctx).CreateInBatches(rows, insertBatchSize).Error
|
||||
}
|
||||
|
||||
@@ -1365,6 +1412,11 @@ func (s *userAccessLogGormStore) BatchInsert(ctx context.Context, logs []analyti
|
||||
if err := s.ensureWritable(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
for i := range logs {
|
||||
if logs[i].ID == 0 {
|
||||
logs[i].ID = idgen.NextUint64ID()
|
||||
}
|
||||
}
|
||||
return s.db.WithContext(ctx).CreateInBatches(logs, insertBatchSize).Error
|
||||
}
|
||||
|
||||
|
||||
@@ -4,9 +4,11 @@
|
||||
package logstore
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -45,6 +47,70 @@ func TestGormBatchInsertAndCount(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// testLogCaptureWriter 捕获 GORM logger 输出(logger.Writer 需实现 Printf)。
|
||||
type testLogCaptureWriter struct {
|
||||
buf *bytes.Buffer
|
||||
}
|
||||
|
||||
func (w testLogCaptureWriter) Write(p []byte) (int, error) { return w.buf.Write(p) }
|
||||
func (w testLogCaptureWriter) Printf(format string, args ...any) {
|
||||
fmt.Fprintf(w.buf, format, args...)
|
||||
}
|
||||
|
||||
// TestGormBatchInsertFillsZeroIDs 回归测试:PG 日志表 id 为 NOT NULL 且无默认值,GORM 对零值
|
||||
// uint64 主键(视为自增)会省略 id 列,导致 PG 插入报 not-null 违例(SQLSTATE 23502)。
|
||||
// 验证 BatchInsert* 落库前为零 ID 行生成雪花 ID:INSERT 语句必须包含 id 列且回填非零、唯一 ID。
|
||||
func TestGormBatchInsertFillsZeroIDs(t *testing.T) {
|
||||
ResetForTest()
|
||||
SetConfigReader(func(_ context.Context, _ string) (string, error) { return "", nil })
|
||||
defer ResetForTest()
|
||||
|
||||
var buf bytes.Buffer
|
||||
dsn := fmt.Sprintf("file:logstore-idtest-%d?mode=memory&cache=shared", atomic.AddInt64(&testGormStoreSeq, 1))
|
||||
db, err := gorm.Open(sqlite.Open(dsn), &gorm.Config{
|
||||
Logger: logger.New(testLogCaptureWriter{&buf}, logger.Config{LogLevel: logger.Info, Colorful: false}),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("open sqlite: %v", err)
|
||||
}
|
||||
if err := db.AutoMigrate(&analyticsmodel.NodeAccessLog{}); err != nil {
|
||||
t.Fatalf("automigrate: %v", err)
|
||||
}
|
||||
s := newGormStore(db)
|
||||
|
||||
now := time.Now()
|
||||
rows := []analyticsmodel.NodeAccessLog{
|
||||
{NodeID: "n1", LoggedAt: now, RemoteAddr: "1.1.1.1", StatusCode: 200},
|
||||
{NodeID: "n1", LoggedAt: now.Add(time.Second), RemoteAddr: "2.2.2.2", StatusCode: 500},
|
||||
}
|
||||
if err := s.BatchInsertNodeAccessLogs(context.Background(), rows); err != nil {
|
||||
t.Fatalf("insert with zero ids: %v", err)
|
||||
}
|
||||
if rows[0].ID == 0 || rows[1].ID == 0 {
|
||||
t.Fatalf("zero ids not filled: %+v %+v", rows[0], rows[1])
|
||||
}
|
||||
if rows[0].ID == rows[1].ID {
|
||||
t.Fatalf("ids not unique: %d == %d", rows[0].ID, rows[1].ID)
|
||||
}
|
||||
// 捕获日志含 CREATE TABLE 等其它语句,仅校验 INSERT 语句的列清单(而非 RETURNING 子句,
|
||||
// 后者无论是否省略列都含 id)。
|
||||
var insertStmt string
|
||||
for _, line := range strings.Split(buf.String(), "\n") {
|
||||
if strings.Contains(line, "INSERT INTO") {
|
||||
insertStmt = line
|
||||
break
|
||||
}
|
||||
}
|
||||
start := strings.Index(insertStmt, "(")
|
||||
end := strings.Index(insertStmt, ") VALUES")
|
||||
if start < 0 || end <= start {
|
||||
t.Fatalf("cannot parse insert statement: %s", insertStmt)
|
||||
}
|
||||
if columns := insertStmt[start+1 : end]; !strings.Contains(columns, "`id`") {
|
||||
t.Fatalf("insert SQL omits id column: %s", insertStmt)
|
||||
}
|
||||
}
|
||||
|
||||
// TestGormNodeAccessLogPagination 验证节点访问日志分页与 CH ListNodeAccessLogs 一致(0-based):
|
||||
// Page=1 size=2 → OFFSET 2;Page=0 视为第 0 页;PageSize<=0 时与 CH 一致不分页(返回全部匹配行)。
|
||||
func TestGormNodeAccessLogPagination(t *testing.T) {
|
||||
@@ -165,6 +231,44 @@ func TestGormNodeAggregatesExcludeEmptyNodeID(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestGormRegionCountsEmptyNodeIDAggregatesAll 回归测试:首页「来源分布」以空 node_id
|
||||
// 表示全节点聚合,RegionCounts 不得拼出 `node_id = ”` 恒空条件(对齐 CH 语义)。
|
||||
func TestGormRegionCountsEmptyNodeIDAggregatesAll(t *testing.T) {
|
||||
ResetForTest()
|
||||
SetConfigReader(func(_ context.Context, _ string) (string, error) { return "", nil })
|
||||
s := newTestGormStore(t)
|
||||
ctx := context.Background()
|
||||
now := time.Now()
|
||||
rows := []analyticsmodel.NodeAccessLog{
|
||||
{ID: 1, NodeID: "n1", LoggedAt: now, RemoteAddr: "1.1.1.1", Region: "CN"},
|
||||
{ID: 2, NodeID: "n2", LoggedAt: now, RemoteAddr: "2.2.2.2", Region: "CN"},
|
||||
{ID: 3, NodeID: "n3", LoggedAt: now, RemoteAddr: "3.3.3.3", Region: "US"},
|
||||
{ID: 4, NodeID: "n4", LoggedAt: now, RemoteAddr: "4.4.4.4", Region: " "},
|
||||
}
|
||||
if err := s.BatchInsertNodeAccessLogs(ctx, rows); err != nil {
|
||||
t.Fatalf("insert: %v", err)
|
||||
}
|
||||
|
||||
all, err := s.RegionCounts(ctx, "", now.Add(-time.Hour), 0)
|
||||
if err != nil {
|
||||
t.Fatalf("region counts (all nodes): %v", err)
|
||||
}
|
||||
if len(all) != 2 {
|
||||
t.Fatalf("all-nodes region counts want 2 regions (empty region excluded), got %+v", all)
|
||||
}
|
||||
if all[0].Region != "CN" || all[0].Count != 2 || all[1].Region != "US" || all[1].Count != 1 {
|
||||
t.Fatalf("all-nodes region counts got %+v, want CN=2 US=1", all)
|
||||
}
|
||||
|
||||
cnOnly, err := s.RegionCounts(ctx, "n1", now.Add(-time.Hour), 0)
|
||||
if err != nil {
|
||||
t.Fatalf("region counts (node): %v", err)
|
||||
}
|
||||
if len(cnOnly) != 1 || cnOnly[0].Region != "CN" || cnOnly[0].Count != 1 {
|
||||
t.Fatalf("node-scoped region counts got %+v, want CN=1", cnOnly)
|
||||
}
|
||||
}
|
||||
|
||||
// testGormStoreSeq 保证每个测试获得独立的共享内存库(cache=shared 下同名 DSN 会复用同一库,
|
||||
// 导致跨测试 id 冲突)。
|
||||
var testGormStoreSeq int64
|
||||
@@ -605,9 +709,9 @@ func TestGormBucketAggregatesFullFieldSet(t *testing.T) {
|
||||
rows := []analyticsmodel.NodeAccessLog{
|
||||
{ID: 1, NodeID: "n1", LoggedAt: base, RemoteAddr: "1.1.1.1", Host: "a.example.com", StatusCode: 200, BytesSent: 100, RequestLength: 10},
|
||||
{ID: 2, NodeID: "n1", LoggedAt: base.Add(time.Minute), RemoteAddr: "1.1.1.1", Host: "a.example.com", StatusCode: 301, BytesSent: 200, RequestLength: 20},
|
||||
{ID: 3, NodeID: "n1", LoggedAt: base.Add(2 * time.Minute), RemoteAddr: "2.2.2.2", Host: "b.example.com", StatusCode: 404, BytesSent: 300, RequestLength: 30},
|
||||
{ID: 3, NodeID: "n1", LoggedAt: base.Add(2 * time.Minute), RemoteAddr: "2.2.2.2", Host: "b.example.com", StatusCode: 400, BytesSent: 300, RequestLength: 30},
|
||||
{ID: 4, NodeID: "n1", LoggedAt: base.Add(3 * time.Minute), RemoteAddr: "", Host: "b.example.com", StatusCode: 500, BytesSent: 400, RequestLength: 40},
|
||||
{ID: 5, NodeID: "n1", LoggedAt: base.Add(4 * time.Minute), RemoteAddr: "3.3.3.3", Host: "c.example.com", StatusCode: 502, BytesSent: 500, RequestLength: 50},
|
||||
{ID: 5, NodeID: "n1", LoggedAt: base.Add(4 * time.Minute), RemoteAddr: "3.3.3.3", Host: "c.example.com", StatusCode: 500, BytesSent: 500, RequestLength: 50},
|
||||
{ID: 6, NodeID: "n2", LoggedAt: base.Add(time.Hour), RemoteAddr: "9.9.9.9", Host: "d.example.com", StatusCode: 200, BytesSent: 999, RequestLength: 99},
|
||||
}
|
||||
if err := s.BatchInsertNodeAccessLogs(ctx, rows); err != nil {
|
||||
@@ -637,6 +741,10 @@ func TestGormBucketAggregatesFullFieldSet(t *testing.T) {
|
||||
if b.ServerErrorCount != 2 {
|
||||
t.Errorf("server_error_count = %d, want 2", b.ServerErrorCount)
|
||||
}
|
||||
if b.Status2xxCount != 1 || b.Status4xxCount != 1 || b.Status5xxCount != 2 {
|
||||
t.Errorf("status class counts = 2xx:%d 4xx:%d 5xx:%d, want 1/1/2 (200/301/400/500/500)",
|
||||
b.Status2xxCount, b.Status4xxCount, b.Status5xxCount)
|
||||
}
|
||||
if b.UniqueIPCount != 3 {
|
||||
t.Errorf("unique_ip_count = %d, want 3 (empty remote_addr excluded)", b.UniqueIPCount)
|
||||
}
|
||||
|
||||
@@ -258,6 +258,9 @@ func buildOpenFlareAccessLogBucketRows(ctx context.Context, query model.OpenFlar
|
||||
SuccessCount: partial.SuccessCount,
|
||||
ClientErrorCount: partial.ClientErrorCount,
|
||||
ServerErrorCount: partial.ServerErrorCount,
|
||||
Status2xxCount: partial.Status2xxCount,
|
||||
Status4xxCount: partial.Status4xxCount,
|
||||
Status5xxCount: partial.Status5xxCount,
|
||||
BytesSent: partial.BytesSent,
|
||||
RequestLength: partial.RequestLength,
|
||||
})
|
||||
|
||||
@@ -10,9 +10,17 @@ const (
|
||||
// OriginErrorPageSupportPath is the SupportFile path for the origin error HTML template.
|
||||
OriginErrorPageSupportPath = "error_pages/origin_error.html.tmpl"
|
||||
|
||||
// OriginErrorPageInternalLocation is the internal nginx location that serves the error body.
|
||||
OriginErrorPageInternalLocation = "/__openflare_origin_error"
|
||||
|
||||
// OriginErrorPageInternalLocation is the named nginx location that serves the error body.
|
||||
// OriginErrorPageInternalLocation is the named nginx location that serves the error body
|
||||
// for the all-methods mode (get_only disabled). Must be a NAMED location (@...), not a
|
||||
// URI internal redirect: error_page URI redirects rewrite the request method to GET, so a
|
||||
// method check inside the location could never distinguish POST/PUT. Named locations keep
|
||||
// the original method and (without the `=` form) the original error status.
|
||||
//
|
||||
// When get_only is enabled this location is NOT emitted: GET-only mode replaces the body
|
||||
// via Lua header/body filters inside the proxy location, so non-GET responses pass through
|
||||
// with their original status and body.
|
||||
OriginErrorPageInternalLocation = "@__openflare_origin_error"
|
||||
defaultOriginErrorPageStatusTag = "500-599"
|
||||
)
|
||||
|
||||
@@ -127,25 +135,83 @@ func renderOriginErrorPageIntercept(cfg ConfigSnapshot) string {
|
||||
if !cfg.OriginErrorPageEnabled {
|
||||
return ""
|
||||
}
|
||||
if _, err := ExpandStatusCodeTags(effectiveOriginErrorPageStatusTags(cfg)); err != nil {
|
||||
codes, err := ExpandStatusCodeTags(effectiveOriginErrorPageStatusTags(cfg))
|
||||
if err != nil || len(codes) == 0 {
|
||||
return ""
|
||||
}
|
||||
if cfg.OriginErrorPageGetOnly {
|
||||
// GET-only mode must NOT use proxy_intercept_errors: interception discards
|
||||
// the upstream error body, so non-GET requests could never receive the
|
||||
// original response (nginx would serve its own default error page instead).
|
||||
// The body is replaced by Lua header/body filters that only fire for GET;
|
||||
// non-GET responses pass through with status, headers and body untouched.
|
||||
return renderOriginErrorPageLuaFilterBlock(codes)
|
||||
}
|
||||
// Intercept at the proxy level for all methods. nginx does not allow
|
||||
// proxy_intercept_errors inside limit_except (only allow/deny are valid
|
||||
// there), so GET-only is enforced in the internal error location's Lua:
|
||||
// non-GET requests exit with the original status and no custom HTML.
|
||||
// there), so the custom HTML is served by the named error location.
|
||||
return " proxy_intercept_errors on;\n"
|
||||
}
|
||||
|
||||
// renderOriginErrorPageServerBits emits server-level error_page + internal location.
|
||||
// Returns empty string when disabled, expand fails, or no codes remain.
|
||||
// renderOriginErrorPageLuaFilterBlock emits the GET-only body replacement inside the
|
||||
// proxy location. header_filter decides whether the response should be replaced and
|
||||
// reads the template once into ngx.ctx; body_filter swaps the upstream body for the
|
||||
// custom HTML and forces end-of-body so remaining upstream chunks are discarded.
|
||||
// Non-GET requests (or statuses outside the configured set) are never touched.
|
||||
func renderOriginErrorPageLuaFilterBlock(codes []int) string {
|
||||
codeList := make([]string, len(codes))
|
||||
for i, code := range codes {
|
||||
codeList[i] = strconv.Itoa(code)
|
||||
}
|
||||
return fmt.Sprintf(` header_filter_by_lua_block {
|
||||
local codes = {%s}
|
||||
local function match(code)
|
||||
for _, c in ipairs(codes) do
|
||||
if c == code then
|
||||
return true
|
||||
end
|
||||
end
|
||||
return false
|
||||
end
|
||||
local status = ngx.status
|
||||
if match(status) and ngx.req.get_method() == "GET" then
|
||||
ngx.header.content_length = nil
|
||||
ngx.header["Content-Type"] = "text/html; charset=utf-8"
|
||||
local f = io.open("%s", "r")
|
||||
local body = f and f:read("*a")
|
||||
if f then
|
||||
f:close()
|
||||
end
|
||||
if not body then
|
||||
body = "<!DOCTYPE html><html><head><meta charset=\"utf-8\"><title>" .. tostring(status) .. "</title></head><body><h1>" .. tostring(status) .. "</h1></body></html>"
|
||||
end
|
||||
body = body:gsub("{{status}}", function() return tostring(status) end)
|
||||
body = body:gsub("{{host}}", function() return ngx.var.host or "" end)
|
||||
ngx.ctx.openflare_error_html = body
|
||||
end
|
||||
}
|
||||
body_filter_by_lua_block {
|
||||
local html = ngx.ctx.openflare_error_html
|
||||
if html then
|
||||
ngx.arg[1] = html
|
||||
ngx.arg[2] = true
|
||||
ngx.ctx.openflare_error_html = nil
|
||||
end
|
||||
}
|
||||
`, strings.Join(codeList, ", "), ErrorPageTmplPlaceholder)
|
||||
}
|
||||
|
||||
// renderOriginErrorPageServerBits emits server-level error_page + named error location
|
||||
// for the all-methods mode. Returns empty string when disabled, expand fails, no codes
|
||||
// remain, or get_only is enabled (GET-only mode replaces the body via Lua filters inside
|
||||
// the proxy location, see renderOriginErrorPageIntercept).
|
||||
//
|
||||
// IMPORTANT: do NOT use `error_page CODE = /uri` (equals without response code).
|
||||
// IMPORTANT: do NOT use `error_page CODE = @name` (equals without response code).
|
||||
// That form adopts the status returned by the error URI; content_by_lua defaults
|
||||
// to 200 and ngx.status is often 0, so clients saw 200 with body "{{status}}"→"0".
|
||||
// Without `=`, nginx keeps the original error status for the internal redirect.
|
||||
// Without `=`, nginx keeps the original error status for the redirect.
|
||||
func renderOriginErrorPageServerBits(cfg ConfigSnapshot) string {
|
||||
if !cfg.OriginErrorPageEnabled {
|
||||
if !cfg.OriginErrorPageEnabled || cfg.OriginErrorPageGetOnly {
|
||||
return ""
|
||||
}
|
||||
codes, err := ExpandStatusCodeTags(effectiveOriginErrorPageStatusTags(cfg))
|
||||
@@ -159,30 +225,25 @@ func renderOriginErrorPageServerBits(cfg ConfigSnapshot) string {
|
||||
var builder strings.Builder
|
||||
// No `=` — preserve original error status (502 stays 502).
|
||||
fmt.Fprintf(&builder, " error_page %s %s;\n", strings.Join(parts, " "), OriginErrorPageInternalLocation)
|
||||
builder.WriteString(renderOriginErrorPageInternalLocation(cfg.OriginErrorPageGetOnly))
|
||||
builder.WriteString(renderOriginErrorPageInternalLocation())
|
||||
return builder.String()
|
||||
}
|
||||
|
||||
func renderOriginErrorPageInternalLocation(getOnly bool) string {
|
||||
// Resolve status from $status (set by error_page internal redirect), then
|
||||
func renderOriginErrorPageInternalLocation() string {
|
||||
// Resolve status from $status (set by error_page redirect), then
|
||||
// upstream_status, then ngx.status. Force ngx.status so the client receives
|
||||
// the real error code. Use function replacers so host/status with `%` are safe.
|
||||
//
|
||||
// Note: fmt.Sprintf is used only for the path placeholders; Lua `%` must be
|
||||
// written as `%%` so Sprintf does not treat them as format verbs.
|
||||
//
|
||||
// When getOnly is true, non-GET that still hit this location (e.g. nginx-local
|
||||
// 502 without upstream body) exit with the original status and no HTML body.
|
||||
getOnlyLua := "false"
|
||||
if getOnly {
|
||||
getOnlyLua = "true"
|
||||
}
|
||||
return fmt.Sprintf(` location = %s {
|
||||
internal;
|
||||
// The location is NAMED (@...), not a URI internal redirect: URI redirects
|
||||
// (location = /uri) rewrite the request method to GET. Named locations keep
|
||||
// the original method and (without `=`) the original error status.
|
||||
return fmt.Sprintf(` location %s {
|
||||
default_type text/html;
|
||||
charset utf-8;
|
||||
content_by_lua_block {
|
||||
local get_only = %s
|
||||
local function resolve_error_status()
|
||||
local code = tonumber(ngx.var.status)
|
||||
if code and code >= 400 then
|
||||
@@ -205,11 +266,6 @@ func renderOriginErrorPageInternalLocation(getOnly bool) string {
|
||||
local code = resolve_error_status()
|
||||
ngx.status = code
|
||||
|
||||
if get_only and ngx.req.get_method() ~= "GET" then
|
||||
-- Non-GET: do not replace with HTML; exit with status only.
|
||||
return ngx.exit(code)
|
||||
end
|
||||
|
||||
local f = io.open("%s", "r")
|
||||
if not f then
|
||||
ngx.header["Content-Type"] = "text/html; charset=utf-8"
|
||||
@@ -227,5 +283,5 @@ func renderOriginErrorPageInternalLocation(getOnly bool) string {
|
||||
ngx.say(body)
|
||||
}
|
||||
}
|
||||
`, OriginErrorPageInternalLocation, getOnlyLua, ErrorPageTmplPlaceholder)
|
||||
`, OriginErrorPageInternalLocation, ErrorPageTmplPlaceholder)
|
||||
}
|
||||
|
||||
@@ -24,22 +24,27 @@ func TestRenderOriginErrorPageEnabled(t *testing.T) {
|
||||
if !strings.Contains(out, "proxy_intercept_errors on") {
|
||||
t.Fatal("missing intercept")
|
||||
}
|
||||
if !strings.Contains(out, "error_page") || !strings.Contains(out, "/__openflare_origin_error") {
|
||||
if !strings.Contains(out, "error_page") || !strings.Contains(out, "@__openflare_origin_error") {
|
||||
t.Fatal("missing error_page")
|
||||
}
|
||||
if !strings.Contains(out, "error_page 500") {
|
||||
t.Fatalf("expected expanded status codes in error_page, got:\n%s", out)
|
||||
}
|
||||
// Must NOT use `error_page … = /uri` (adopts error-URI status → often 200).
|
||||
// `location = /path` is unrelated and expected.
|
||||
// Must NOT use `error_page … = @name` (adopts error-URI status → often 200).
|
||||
for _, line := range strings.Split(out, "\n") {
|
||||
trimmed := strings.TrimSpace(line)
|
||||
if strings.HasPrefix(trimmed, "error_page ") && strings.Contains(trimmed, " = ") {
|
||||
t.Fatalf("error_page must not use '=' form, got: %s", trimmed)
|
||||
}
|
||||
}
|
||||
if !strings.Contains(out, "error_page ") || !strings.Contains(out, " /__openflare_origin_error;") {
|
||||
t.Fatal("error_page must redirect to internal location without '='")
|
||||
if !strings.Contains(out, "error_page ") || !strings.Contains(out, " @__openflare_origin_error;") {
|
||||
t.Fatal("error_page must redirect to the named error location without '='")
|
||||
}
|
||||
if !strings.Contains(out, "location @__openflare_origin_error {") {
|
||||
t.Fatal("error location must be a named location (@...) that preserves the request method")
|
||||
}
|
||||
if strings.Contains(out, "location = /__openflare_origin_error") {
|
||||
t.Fatal("error location must NOT be a URI internal redirect (error_page URI redirects rewrite the method to GET, breaking the get_only gate)")
|
||||
}
|
||||
if !strings.Contains(out, "resolve_error_status") || !strings.Contains(out, "ngx.status = code") {
|
||||
t.Fatal("internal location must resolve and set ngx.status to the original error code")
|
||||
@@ -85,23 +90,46 @@ func TestRenderOriginErrorPageGetOnly(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !strings.Contains(out, "proxy_intercept_errors on") {
|
||||
t.Fatal("missing intercept on")
|
||||
// Regression: GET-only must NOT intercept at the proxy level.
|
||||
// proxy_intercept_errors discards the upstream error body, so non-GET requests
|
||||
// would receive nginx's own default error page instead of the original
|
||||
// response (this was the reported bug: POST 503 returned OpenResty's page).
|
||||
if strings.Contains(out, "proxy_intercept_errors") {
|
||||
t.Fatal("get_only must not emit proxy_intercept_errors (it discards the upstream body for non-GET)")
|
||||
}
|
||||
// The body replacement must happen in Lua filters that only fire for GET.
|
||||
if !strings.Contains(out, "header_filter_by_lua_block") {
|
||||
t.Fatal("get_only must emit header_filter_by_lua_block inside the proxy location")
|
||||
}
|
||||
if !strings.Contains(out, "body_filter_by_lua_block") {
|
||||
t.Fatal("get_only must emit body_filter_by_lua_block inside the proxy location")
|
||||
}
|
||||
if !strings.Contains(out, `ngx.req.get_method() == "GET"`) {
|
||||
t.Fatal("Lua filter must replace the body only for GET requests")
|
||||
}
|
||||
if !strings.Contains(out, `ngx.ctx.openflare_error_html`) {
|
||||
t.Fatal("Lua filter must stash the error HTML in ngx.ctx for the body filter")
|
||||
}
|
||||
if !strings.Contains(out, `local codes = {500`) {
|
||||
t.Fatal("Lua filter must carry the expanded status codes")
|
||||
}
|
||||
if !strings.Contains(out, ErrorPageTmplPlaceholder) {
|
||||
t.Fatal("missing error page template placeholder")
|
||||
}
|
||||
// No error_page / named location machinery in GET-only mode.
|
||||
if strings.Contains(out, "error_page") {
|
||||
t.Fatal("get_only must not emit error_page (named-location path can only serve HTML or an empty status, never the original body)")
|
||||
}
|
||||
if strings.Contains(out, "@__openflare_origin_error") {
|
||||
t.Fatal("get_only must not emit the named error location")
|
||||
}
|
||||
// nginx rejects proxy_intercept_errors inside limit_except (only allow/deny
|
||||
// are valid there), which made the generated config fail `openresty -t` and
|
||||
// caused apply rollback. GET-only must rely on the internal location's Lua.
|
||||
// are valid there); GET-only must rely on Lua filters instead.
|
||||
if strings.Contains(out, "limit_except") {
|
||||
t.Fatal("get_only must not emit limit_except (proxy_intercept_errors is not allowed there)")
|
||||
}
|
||||
if strings.Contains(out, "proxy_intercept_errors off") {
|
||||
t.Fatal("get_only must not emit proxy_intercept_errors off")
|
||||
}
|
||||
if !strings.Contains(out, `get_only = true`) {
|
||||
t.Fatal("internal location must set get_only = true")
|
||||
}
|
||||
if !strings.Contains(out, `ngx.req.get_method() ~= "GET"`) {
|
||||
t.Fatal("internal location must skip HTML for non-GET")
|
||||
if strings.Contains(out, "location = /__openflare_origin_error") {
|
||||
t.Fatal("must not use URI internal redirect (rewrites method to GET, breaking the GET gate)")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -121,8 +149,8 @@ func TestRenderOriginErrorPageDisabled(t *testing.T) {
|
||||
if strings.Contains(out, "proxy_intercept_errors") {
|
||||
t.Fatal("should not intercept when disabled")
|
||||
}
|
||||
if strings.Contains(out, "/__openflare_origin_error") {
|
||||
t.Fatal("should not emit internal error location when disabled")
|
||||
if strings.Contains(out, "@__openflare_origin_error") {
|
||||
t.Fatal("should not emit error location when disabled")
|
||||
}
|
||||
res, err := Render(doc, nil)
|
||||
if err != nil {
|
||||
@@ -199,7 +227,7 @@ func TestRenderOriginErrorPageCustomHTMLInSupportFile(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !strings.Contains(out, "error_page 502 /__openflare_origin_error;") {
|
||||
if !strings.Contains(out, "error_page 502 @__openflare_origin_error;") {
|
||||
t.Fatalf("expected single 502 error_page without '=', got:\n%s", out)
|
||||
}
|
||||
}
|
||||
@@ -227,7 +255,7 @@ func TestRenderOriginErrorPageSkipsPagesRoutes(t *testing.T) {
|
||||
if strings.Contains(out, "proxy_intercept_errors") {
|
||||
t.Fatal("pages routes must not get proxy_intercept_errors")
|
||||
}
|
||||
if strings.Contains(out, "/__openflare_origin_error") {
|
||||
if strings.Contains(out, "@__openflare_origin_error") {
|
||||
t.Fatal("pages routes must not get origin error location")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user