Compare commits

...

4 Commits

Author SHA1 Message Date
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
10 changed files with 408 additions and 146 deletions
-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
+60 -105
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`.
> After the first login with the `admin` user, you must change the default password `12345678`.
>
> 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.
> 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
+4
View File
@@ -18,6 +18,10 @@ sidebar: false
## [Unreleased]
### 🛠 修复
- 修复源站错误页「仅针对 GET 请求」未生效:`error_page` 内部重定向会把请求方法改写成 GET,导致内部 Lua 无法识别 POST/PUT 等原始方法、仍返回自定义错误页;现改为命名 location(`@__openflare_origin_error`)承载错误页,保留原始请求方法与错误状态码,非 GET 请求不再返回自定义错误页。
- 修复 PostgreSQL 作为日志库时节点访问日志/可观测指标/用户访问日志批量写入失败:GORM 对零值 `uint64` 主键会省略 `id` 列,而 PG 日志表 `id` 无默认值,导致持续报「null value in column id violates not-null constraint」;现于落库前为零 ID 行生成雪花 ID(与 ClickHouse 写入路径一致),并新增回归测试覆盖六张日志表。
## [v3.5.1] - 2026-08-09
### 新增
-1
View File
@@ -55,7 +55,6 @@ import {
LogOut,
Settings,
ShieldCheck,
Terminal,
UserRound,
} from 'lucide-react';
@@ -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
}
@@ -1000,6 +1008,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 +1213,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 +1268,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 +1323,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 +1393,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) {
+15 -10
View File
@@ -10,8 +10,12 @@ 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.
// Must be a NAMED location (@...), not a URI internal redirect: error_page URI redirects
// rewrite the request method to GET, so the get_only Lua check (ngx.req.get_method() ~= "GET")
// would never fire and POST/PUT would still receive the custom HTML. Named locations preserve
// the original request method and the original error status (without the `=` form).
OriginErrorPageInternalLocation = "@__openflare_origin_error"
defaultOriginErrorPageStatusTag = "500-599"
)
@@ -137,13 +141,13 @@ func renderOriginErrorPageIntercept(cfg ConfigSnapshot) string {
return " proxy_intercept_errors on;\n"
}
// renderOriginErrorPageServerBits emits server-level error_page + internal location.
// renderOriginErrorPageServerBits emits server-level error_page + named error location.
// Returns empty string when disabled, expand fails, or no codes remain.
//
// 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 {
return ""
@@ -164,21 +168,22 @@ func renderOriginErrorPageServerBits(cfg ConfigSnapshot) string {
}
func renderOriginErrorPageInternalLocation(getOnly bool) string {
// Resolve status from $status (set by error_page internal redirect), then
// 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.
// The location is NAMED (@...), not a URI internal redirect: URI redirects
// (location = /uri) rewrite the request method to GET, which defeats the
// get_only gate below. Named locations keep the original method, so a POST
// that reaches this location exits with the original status and no HTML body.
getOnlyLua := "false"
if getOnly {
getOnlyLua = "true"
}
return fmt.Sprintf(` location = %s {
internal;
return fmt.Sprintf(` location %s {
default_type text/html;
charset utf-8;
content_by_lua_block {
+24 -9
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")
@@ -103,6 +108,16 @@ func TestRenderOriginErrorPageGetOnly(t *testing.T) {
if !strings.Contains(out, `ngx.req.get_method() ~= "GET"`) {
t.Fatal("internal location must skip HTML for non-GET")
}
// Regression: error_page must target the NAMED location. URI internal redirects
// (location = /uri) rewrite the request method to GET, so ngx.req.get_method()
// would always return "GET" and the get_only gate would never skip HTML for
// POST/PUT. Named locations preserve the original method.
if !strings.Contains(out, "error_page 500") || !strings.Contains(out, " @__openflare_origin_error;") {
t.Fatalf("error_page must target the named location, got:\n%s", out)
}
if strings.Contains(out, "location = /__openflare_origin_error") {
t.Fatal("get_only must not use URI internal redirect (rewrites method to GET, breaking the gate)")
}
}
func TestRenderOriginErrorPageDisabled(t *testing.T) {
@@ -121,8 +136,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 +214,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 +242,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")
}
}