Compare commits

...

10 Commits

Author SHA1 Message Date
ryan f960511cc0 chore(release): v3.5.3
### 新增
- 访问日志「日志明细」支持按 HTTP 状态码筛选,可直接输入任意状态码。
- 访问日志「日志明细」支持自定义时间范围筛选,可按起止时间检索日志。
- 首页看板改版:24 小时请求趋势拆分展示请求总量与 2xx/4xx/5xx 状态码类请求量并独占一行;移除宿主机磁盘指标,24 小时容量趋势(CPU/内存)并入业务流量卡片展示。

### 🛠 修复
- 修复首页「来源分布」卡片在 PostgreSQL/SQLite 日志库下无数据的问题。
- 修复源站错误页「仅针对 GET 请求」未真正透传非 GET 响应的问题:POST/PUT 等非 GET 请求现可完整看到源站原始报错内容。
2026-08-13 11:44:08 +08:00
ryan 465440fa5b fix(access-logs): 修复状态码自定义 2026-08-13 11:33:12 +08:00
ryan a4dd5ca9e1 feat(dashboard): 首页请求趋势拆分状态码并合并容量到业务流量
- 24 小时请求趋势拆分展示请求总量与 200/400/500 状态码请求量,独占一行;
  时间桶聚合新增 status_200/400/500_count(CH countIf、PG FILTER),
  请求趋势改为基于原始桶聚合(小时 rollup 无状态码口径)
- 首页移除宿主机磁盘指标,容量趋势(CPU/内存)并入业务流量卡片展示
- 压缩协议 traffic_24h 扩展为 7 元组,前端归一化同步更新
2026-08-13 11:10:37 +08:00
ryan a9e4237bbf feat(access-logs): 状态码支持手动输入,新增时间范围筛选
- 状态码筛选支持预设快捷选项 + 手动输入任意 100-599 状态码(数字校验)
- 新增时间范围筛选:shadcn 日期+时间选择器(Popover+Calendar+时分 Select),
  起止时间以 RFC3339 成对传入,后端校验格式与先后关系,非法值返回 400
- 默认显示来源 IP/访问域名/状态码,节点 ID/请求路径/时间范围折叠进「更多筛选」
2026-08-13 10:27:40 +08:00
ryan 75d1fcf345 feat(access-logs): 日志明细支持按状态码筛选并折叠次要搜索项,修复首页来源分布无数据
- 修复 PostgreSQL/SQLite 日志库下首页「来源分布」卡片无数据:RegionCounts 对空
  节点 ID 误拼 node_id = '' 恒空条件,改为空节点 ID 表示全节点聚合(对齐 CH 语义),
  并过滤空白归属地
- /access-logs?tab=list 新增状态码筛选:状态码下拉含常用 2xx/3xx/4xx/5xx 选项,
  校验 100-599,非法值返回 400;ClickHouse 与 PostgreSQL/SQLite 日志库均支持
- 搜索框折叠:默认仅显示来源 IP 与状态码,节点 ID/访问域名/请求路径折叠进
  「更多筛选」
2026-08-13 09:59:32 +08:00
ryan f1577bf092 fix(openresty): 修复源站错误页「仅针对 GET 请求」覆盖非 GET 原始报错数据
proxy_intercept_errors 会在 Lua 判断前丢弃源站错误响应体,POST/PUT 等
请求收到 503 时被 OpenResty 自带错误页覆盖原始报错数据。现改为在代理
location 内用 Lua header/body 过滤器仅对 GET 请求替换错误页,非 GET
请求完整透传源站原始状态码与响应体;非仅 GET 模式继续使用命名 location
承载错误页。
2026-08-09 19:38:44 +08:00
ryan 01ed2c5e36 chore(release): v3.5.2
修复几个遗漏bug
2026-08-09 14:08:04 +08:00
ryan 80696c12fa fix: lint 2026-08-09 13:47:40 +08:00
ryan 3d4d99081e fix(log): PG 日志库批量写入为零 ID 行生成雪花 ID
PostgreSQL 日志表 id 为 NOT NULL 且无默认值,而 GORM 将零值 uint64
主键视为自增并省略 id 列,导致 node access log / 可观测指标等批量
落库持续报 "null value in column id violates not-null constraint"。
在 BatchInsert* 落库前为零 ID 行生成雪花 ID(与 ClickHouse 写入路径
一致),并新增单元回归与 PG 集成回归测试覆盖六张日志表。
2026-08-09 13:47:13 +08:00
ryan 0639855653 fix(openresty): 修复源站错误页「仅针对 GET 请求」未生效
error_page 的 URI 内部重定向会把请求方法改写成 GET,导致内部
Lua 中 ngx.req.get_method() 恒为 GET,get_only 判断永不命中,
POST/PUT 等请求仍返回自定义错误页。

改为命名 location(@__openflare_origin_error)承载错误页:
命名 location 保留原始请求方法与原始错误状态码,非 GET 请求
直接以原状态码退出、不再注入自定义 HTML。附带回归断言,禁止
回退到 URI 内部重定向形式。
2026-08-09 13:42:38 +08:00
38 changed files with 1141 additions and 385 deletions
+7 -11
View File
@@ -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 变化产生的冗余令牌。
-19
View File
@@ -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
View File
@@ -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
![OpenFlare dashboard overview](./docs/assets/readme/dashboard-overview.png)
### Access Logs
![OpenFlare version release](./docs/assets/readme/domain_overview.png)
### WAF Protection
![OpenFlare version release](./docs/assets/readme/waf.png)
## 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
![OpenFlare dashboard overview](./docs/assets/readme/dashboard-overview.png)
### Node Details
![OpenFlare node detail](./docs/assets/readme/node-detail.png)
### Proxy Configuration
![OpenFlare version release](./docs/assets/readme/proxy-route-detail.png)
## 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
View File
@@ -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
View File
@@ -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": "页码",
+1 -1
View File
@@ -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
View File
@@ -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
View File
@@ -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;
+11
View File
@@ -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>
);
}
@@ -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),
},
]}
/>
+9 -11
View File
@@ -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} />
</>
)}
-1
View File
@@ -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,
);
+15 -1
View File
@@ -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];
+1 -1
View File
@@ -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 {
+10 -8
View File
@@ -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"`
+13 -8
View File
@@ -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
)`
+53 -1
View File
@@ -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,
})
+85 -29
View File
@@ -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)
}
+49 -21
View File
@@ -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")
}
}