mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-28 05:46:36 +08:00
Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 01ed2c5e36 | |||
| 80696c12fa | |||
| 3d4d99081e | |||
| 0639855653 |
@@ -1,19 +0,0 @@
|
||||
root = true
|
||||
|
||||
[*]
|
||||
indent_style = space
|
||||
indent_size = 4
|
||||
charset = utf-8
|
||||
end_of_line = lf
|
||||
trim_trailing_whitespace = true
|
||||
insert_final_newline = true
|
||||
|
||||
[*.{json,yml,yaml}]
|
||||
indent_size = 2
|
||||
|
||||
[*.md]
|
||||
insert_final_newline = false
|
||||
trim_trailing_whitespace = false
|
||||
|
||||
[*.{js,ts,css,html,jsx,tsx,vue}]
|
||||
indent_size = 2
|
||||
+61
-106
@@ -2,9 +2,9 @@
|
||||
|
||||
# OpenFlare
|
||||
|
||||
**[English](./README.en.md) | [📖 中文](./README.md)**
|
||||
**[📖 中文](./README.md) | [English](./README.en.md)**
|
||||
|
||||
OpenFlare is an open-source CDN orchestration and edge security platform. It supports reverse proxies, centralized configuration synchronization, secure intranet penetration (Tunnels), dynamic WAF protection, and anti-CC challenges.
|
||||
OpenFlare is an open-source CDN orchestration and edge security platform. It supports reverse proxy, centralized configuration synchronization, in-network tunneling (Tunnels), dynamic WAF protection, and CC defense challenges.
|
||||
|
||||
</div>
|
||||
|
||||
@@ -21,37 +21,67 @@ OpenFlare is an open-source CDN orchestration and edge security platform. It sup
|
||||
</p>
|
||||
|
||||
> [!WARNING]
|
||||
> After logging in for the first time with the `root` user, make sure to change the default password `123456`.
|
||||
>
|
||||
> The BETA version is a temporary product for the development and testing phase. It may contain unknown issues and should not be used in production environments.
|
||||
> After the first login with the `admin` user, you must change the default password `12345678`.
|
||||
>
|
||||
> The BETA version is a temporary product in the development and testing stage and may have unknown issues. It should not be used in production environments.
|
||||
|
||||
## Documentation
|
||||
|
||||
**https://open-flare.pages.dev**
|
||||
|
||||
Quick links:
|
||||
Common entry points:
|
||||
|
||||
* [Quick Start](https://open-flare.pages.dev/en/guide/quick-start)
|
||||
* [Deployment Guide](https://open-flare.pages.dev/en/deployment/deployment)
|
||||
* [Quick Start](https://open-flare.pages.dev/guide/quick-start)
|
||||
* [Deployment Guide](https://open-flare.pages.dev/deployment/deployment)
|
||||
* [Configuration Reference](https://open-flare.pages.dev/reference/configuration)
|
||||
* [System Design](https://open-flare.pages.dev/design/)
|
||||
|
||||
## Core Features
|
||||
## Core Capabilities
|
||||
|
||||
* **Reverse Proxy Management**: Website rules as the aggregation boundary, supporting multi-domain binding and multi-upstream load balancing with unified management of all OpenResty node configurations.
|
||||
* **Immutable Config Version Control**: Full-snapshot publish model based on version numbers (`YYYYMMDD-NNN`), with pre-publish diff preview, a single globally active version, and one-click sub-second rollback.
|
||||
* **Secure Intranet Penetration (Tunnels)**: An open-source alternative to Cloudflare Tunnels. Securely expose local intranet Web services to the public network via Relay and OpenFlared clients — no public IP or open inbound ports required.
|
||||
* **Edge WAF Safety Protection**: Provides global and custom rule groups, supporting manual/automatic/subscription IP groups, MaxMind GeoIP country-level access control, Checksum-based differential IP group sync (no Nginx reload), and custom block responses.
|
||||
* **Anti-CC & Human-Machine Challenge (PoW)**: Built-in high-performance client-side cryptographic Proof of Work challenges (similar to Turnstile) to block and intercept botnets and scrapers at the gateway edge in seconds.
|
||||
* **Pages Static Hosting**: Upload pre-built ZIP packages directly; edge Agents pull and serve them via local OpenResty, with SPA Fallback and built-in API reverse proxy configuration.
|
||||
* **Automated TLS Certificate Management**: Supports dynamic certificate upload, automatic multi-domain certificate matching and binding, and ACME-based automatic issuance and renewal via Let's Encrypt.
|
||||
* **Uptime Kuma Monitoring Sync**: Integrates with Uptime Kuma to automatically sync the monitoring site list using differential updates, providing real-time awareness of node availability and service health.
|
||||
* **SSO Single Sign-On**: Supports GitHub OAuth and standard OIDC protocol for seamless integration with enterprise identity providers.
|
||||
* **Unified Observability**: Aggregates node request metrics, real-time access log details, host/Nginx resource snapshots, health events, and a re-upload buffer for network fluctuations.
|
||||
* **Reverse Proxy Configuration Management**: Uses website rules as the aggregation boundary, supports multi-domain binding and multi-upstream load balancing, and centrally manages reverse proxy configurations for all OpenResty nodes.
|
||||
* **Secure In-Network Tunneling (Tunnels)**: Open-source version of Cloudflare Tunnels. No public IP or exposed inbound ports are required. Securely reverse-proxy internal web services to the public internet through Relay relay nodes and OpenFlared clients.
|
||||
* **Edge WAF Security Protection**: Provides global and custom rule groups, supports manual/auto/subscription-type IP groups, MaxMind GeoIP national-level geographic access control, IP group member Checksum differential synchronization (no Nginx reload required), and custom blocking responses.
|
||||
* **CC Defense and Human-Computer Challenge (PoW)**: Built-in high-performance client-side cryptography Proof of Work challenge (similar to Turnstile). Secures high-speed interception and blocking of zombie networks and crawlers at the gateway edge.
|
||||
* **Pages Static Hosting**: Supports uploading or synchronizing pre-built artifacts from restricted Remote URLs or public GitHub Release assets. GitHub latest can be checked periodically and optionally auto-published. All sources are unified to generate immutable deployments, pulled by the edge Agent and served locally by OpenResty, supporting rollbacks, SPA Fallback, and API reverse proxy.
|
||||
* **TLS Certificate Automation**: Supports dynamic certificate uploads, automatic multi-domain certificate matching and binding, and automatic issuance and renewal of certificates from Let's Encrypt via the ACME protocol.
|
||||
* **Uptime Kuma Monitoring Synchronization**: Integrated with Uptime Kuma to automatically perform differential synchronization of monitoring site lists, real-time awareness of node availability and service status.
|
||||
* **SSO Single Sign-On**: Supports GitHub OAuth and standard OIDC protocol for seamless integration with enterprise identity providers to achieve unified login.
|
||||
* **Unified Observability**: Aggregates node request metrics, real-time access log details, host and Nginx resource snapshots, health events, and network fluctuation replenishment buffers.
|
||||
|
||||
## Interface Preview
|
||||
|
||||
### Dashboard Overview
|
||||
|
||||

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

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

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

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

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

|
||||
|
||||
## License
|
||||
|
||||
This project is licensed under [Apache License 2.0](./LICENSE).
|
||||
This project is licensed under the [Apache License 2.0](./LICENSE).
|
||||
|
||||
## Star History
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
### 新增
|
||||
|
||||
@@ -55,7 +55,6 @@ import {
|
||||
LogOut,
|
||||
Settings,
|
||||
ShieldCheck,
|
||||
Terminal,
|
||||
UserRound,
|
||||
} from 'lucide-react';
|
||||
|
||||
|
||||
@@ -4559,4 +4559,4 @@
|
||||
}
|
||||
}
|
||||
}
|
||||
]
|
||||
]
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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,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")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user