mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-28 05:46:36 +08:00
refactor(edge): Server 协议统一与 wsclient 收敛(Phase 3 Batch 3)
- openflare/agent/relay/flared 协议类型改为 pkg/protocol 别名,删除重复 struct - 新增 edge/wsclient Preset 表,三组件 wsclient 改为薄包装 - 更新设计文档、计划与 changelog
This commit is contained in:
@@ -21,6 +21,8 @@ sidebar: false
|
||||
- 边缘组件运行时去重:新增 `internal/apps/edge/` 共享包(`updater`、`httpclient`、`nodeip`、`logging`、`heartbeat/autoupdate`、`runner`),Agent/Relay/Flared 三组件改为薄包装委托,删除约 1100 行重复自更新与 HTTP 传输层代码。
|
||||
- 边缘运行时 Phase 3 Batch 1:抽取 `edge/config/duration`、`edge/observability/linux`、`edge/heartbeat/loop`,统一 MillisecondDuration、Linux 指标采集与 relay/flared 心跳循环。
|
||||
- 边缘运行时 Phase 3 Batch 2:Agent 心跳周期下沉至 `heartbeat/cycle.go`;`pkg/protocol/agent.go` 统一 Agent 客户端协议类型。
|
||||
- 边缘运行时 Phase 3 Batch 3(T6):Server 侧 `openflare/agent`、`relay`、`flared` 协议类型统一为 `pkg/protocol` 别名,消除与客户端的重复定义。
|
||||
- 边缘运行时 Phase 3 Batch 3(T7):新增 `edge/wsclient`,Agent/Relay/Flared 三组件 `wsclient/` 改为 Preset 薄包装,删除约 150 行重复 WebSocket 传输层代码。
|
||||
- 合并并简化仓库结构:将 `openflare-server` 单体目录下的所有文件/目录提升至仓库根目录(去除了 `openflare-server` 嵌套层级),保留 `.github` 目录不变;统一配置 `docker-compose.yaml` 及所有 Dockerfile 的构建上下文为根目录。
|
||||
- 调整子项目结构与包路径:将 `agent`、`relay` 和 `flared` 子项目从 `internal/` 移动至 `internal/apps/`(分别为 `internal/apps/agent`、`internal/apps/relay` 和 `internal/apps/flared`),并递归更新了所有涉及的 Go 导入路径(如 `github.com/Rain-kl/Wavelet/internal/apps/agent` 等)。
|
||||
- 调整编译产物输出名称与 Makefile:
|
||||
|
||||
@@ -34,6 +34,7 @@ internal/apps/edge/
|
||||
├── logging/ # Setup、ParseLevel
|
||||
├── nodeip/ # Detect、DetectLocal(可注入 LookupOutboundIP)
|
||||
├── httpclient/ # 基础 HTTP 客户端(鉴权头可配置)
|
||||
├── wsclient/ # 基础 WebSocket 客户端(组件 Preset 驱动 HeaderKey + WSPath)
|
||||
├── updater/ # GitHub Release 自更新 + 二进制替换重启
|
||||
├── heartbeat/ # TryAutoUpdate 统一入口
|
||||
└── runner/ # WS 重连循环、SleepContext
|
||||
@@ -49,7 +50,7 @@ internal/apps/edge/
|
||||
| Relay | `frps/`、`observability/` |
|
||||
| Flared | `frpc/`、`sync/`(tunnel) |
|
||||
|
||||
各组件 `updater/`、`httpclient/` 变为类型别名 + `New()` 工厂函数。
|
||||
各组件 `updater/`、`httpclient/`、`wsclient/` 变为类型别名 + `New()` 工厂函数。
|
||||
|
||||
---
|
||||
|
||||
@@ -72,6 +73,14 @@ edgehttp.New(baseURL, token, timeout, "X-Agent-Token") // Agent/Relay
|
||||
edgehttp.New(baseURL, token, timeout, "X-Tunnel-Token") // Flared
|
||||
```
|
||||
|
||||
### WebSocket 客户端
|
||||
|
||||
```go
|
||||
edgews.New(edgews.PresetAgent, baseURL, token, timeout) // HeaderKey=X-Agent-Token, /api/v1/agent/ws
|
||||
edgews.New(edgews.PresetRelay, baseURL, token, timeout) // HeaderKey=X-Agent-Token, /api/v1/relay/ws
|
||||
edgews.New(edgews.PresetFlared, baseURL, token, timeout) // HeaderKey=X-Tunnel-Token, /api/v1/tunnel/ws
|
||||
```
|
||||
|
||||
### 节点 IP 探测
|
||||
|
||||
```go
|
||||
@@ -104,12 +113,15 @@ nodeip.Detect() // outbound → local 回退
|
||||
- [x] `heartbeat/cycle.go` — Agent HTTP 心跳周期从 runner 下沉(payload 构建、同步、自动更新)
|
||||
- [x] `pkg/protocol/agent.go` — Agent 客户端协议类型迁入公共包,`internal/apps/agent/protocol` 保留别名 re-export
|
||||
|
||||
## 已完成(Phase 3 Batch 3)
|
||||
|
||||
- [x] `edge/wsclient` — Agent/Relay/Flared WebSocket 传输层收敛(Preset 配置表 + `AgentConnection` 适配 `protocol.WebSocketConnection`)
|
||||
- [x] Server 侧协议统一 — `internal/apps/openflare/{agent,relay,flared}` 心跳/观测类型改为 `pkg/protocol` 别名
|
||||
|
||||
## 可选后续
|
||||
|
||||
| 项 | 说明 |
|
||||
| --- | --- |
|
||||
| Server 侧协议统一 | 评估 `internal/apps/openflare/agent` 与 `pkg/protocol` 类型去重 |
|
||||
| `wsclient` 薄包装收敛 | relay/flared/agent wsclient 配置表化 |
|
||||
|
||||
---
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
# 边缘运行时 Phase 3 — 任务拆解
|
||||
|
||||
> **状态**:Batch 1 + Batch 2 已完成(2026-06-19)
|
||||
> **状态**:Batch 3 已完成(2026-06-19)
|
||||
> **前置**:[边缘运行时重构设计](../design/edge-runtime-refactor.md) Phase 0–2 已完成
|
||||
|
||||
---
|
||||
@@ -45,9 +45,31 @@ Batch 2(串行,依赖 Batch 1 或需独立评审)
|
||||
|
||||
---
|
||||
|
||||
## 验收标准(Batch 1)
|
||||
## Batch 3 — 并行任务(已完成)
|
||||
|
||||
| ID | 任务 | 修改范围 | 状态 |
|
||||
| --- | --- | --- | --- |
|
||||
| **T6** | Server 侧协议统一 | `pkg/protocol/` + `internal/apps/openflare/{agent,relay,flared}/` | ✅ 子代理 F |
|
||||
| **T7** | `wsclient` 收敛 | `edge/wsclient/` + 三组件 `wsclient/` 薄包装 | ✅ 子代理 G |
|
||||
|
||||
### T6 要点
|
||||
|
||||
- `openflare/agent` 的 `NodePayload`、观测类型、WAF 类型改为 `pkg/protocol` 别名
|
||||
- 保留 Server 专有响应:`RegistrationResponse`(`access_token`)、`HeartbeatResponse`(含 `*model.OpenFlareNode`)
|
||||
- `openflare/relay`、`openflare/flared` 心跳/配置载荷改为 `pkg/protocol` 别名
|
||||
|
||||
### T7 要点
|
||||
|
||||
- 新增 `internal/apps/edge/wsclient/`,配置表驱动(HeaderKey + WSPath)
|
||||
- Agent 保留 `SendStatus` + `protocol.WebSocketConnection` 适配
|
||||
- relay/flared/agent 的 `wsclient/` 仅保留 `New()` 工厂
|
||||
|
||||
---
|
||||
|
||||
## 验收标准
|
||||
|
||||
```bash
|
||||
go build ./cmd/agent ./cmd/relay ./cmd/flared
|
||||
go test ./internal/apps/edge/... ./internal/apps/agent/... ./internal/apps/relay/... ./internal/apps/flared/... -count=1
|
||||
go test ./internal/apps/edge/... ./internal/apps/agent/... ./internal/apps/relay/... ./internal/apps/flared/... ./internal/apps/openflare/... ./pkg/protocol/... -count=1
|
||||
make code-check
|
||||
```
|
||||
@@ -5,78 +5,31 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/agent/protocol"
|
||||
shared "github.com/Rain-kl/Wavelet/pkg/wsclient"
|
||||
edgews "github.com/Rain-kl/Wavelet/internal/apps/edge/wsclient"
|
||||
)
|
||||
|
||||
type WSMessage = shared.WSMessage
|
||||
type MessageHandler = shared.MessageHandler
|
||||
type WSMessage = edgews.WSMessage
|
||||
type MessageHandler = edgews.MessageHandler
|
||||
type Connection = edgews.AgentConnection
|
||||
|
||||
type Client struct {
|
||||
sharedClient *shared.Client
|
||||
inner *edgews.Client
|
||||
}
|
||||
|
||||
type Connection struct {
|
||||
sharedConn *shared.Connection
|
||||
}
|
||||
|
||||
func New(baseURL string, token string, timeout time.Duration) *Client {
|
||||
func New(baseURL, token string, timeout time.Duration) *Client {
|
||||
return &Client{
|
||||
sharedClient: shared.New(shared.Config{
|
||||
BaseURL: baseURL,
|
||||
Token: token,
|
||||
Timeout: timeout,
|
||||
HeaderKey: "X-Agent-Token",
|
||||
WSPath: "/api/v1/agent/ws",
|
||||
}),
|
||||
inner: edgews.New(edgews.PresetAgent, baseURL, token, timeout),
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Client) SetToken(token string) {
|
||||
c.sharedClient.SetToken(token)
|
||||
c.inner.SetToken(token)
|
||||
}
|
||||
|
||||
func (c *Client) URL() string {
|
||||
return c.sharedClient.URL()
|
||||
return c.inner.URL()
|
||||
}
|
||||
|
||||
func (c *Client) Connect(ctx context.Context) (protocol.WebSocketConnection, error) {
|
||||
conn, err := c.sharedClient.Connect(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &Connection{sharedConn: conn}, nil
|
||||
}
|
||||
|
||||
func (conn *Connection) URL() string {
|
||||
if conn == nil || conn.sharedConn == nil {
|
||||
return ""
|
||||
}
|
||||
return conn.sharedConn.URL
|
||||
}
|
||||
|
||||
func (conn *Connection) SendStatus(payload protocol.NodePayload) error {
|
||||
return conn.sharedConn.SendMessage(protocol.WSMessageTypeStatus, payload)
|
||||
}
|
||||
|
||||
func (conn *Connection) SendPong() error {
|
||||
return conn.sharedConn.SendMessage(protocol.WSMessageTypePong, nil)
|
||||
}
|
||||
|
||||
func (conn *Connection) Receive() (protocol.WSMessage, error) {
|
||||
var message protocol.WSMessage
|
||||
if err := conn.sharedConn.Receive(&message); err != nil {
|
||||
return message, err
|
||||
}
|
||||
return message, nil
|
||||
}
|
||||
|
||||
func (conn *Connection) RunReceiveLoop(ctx context.Context, handler shared.MessageHandler) error {
|
||||
return conn.sharedConn.RunReceiveLoop(ctx, handler)
|
||||
}
|
||||
|
||||
func (conn *Connection) Close() error {
|
||||
if conn == nil || conn.sharedConn == nil {
|
||||
return nil
|
||||
}
|
||||
return conn.sharedConn.Close()
|
||||
}
|
||||
return c.inner.ConnectAgent(ctx)
|
||||
}
|
||||
@@ -0,0 +1,130 @@
|
||||
package wsclient
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
pkgprotocol "github.com/Rain-kl/Wavelet/pkg/protocol"
|
||||
shared "github.com/Rain-kl/Wavelet/pkg/wsclient"
|
||||
)
|
||||
|
||||
type WSMessage = shared.WSMessage
|
||||
type MessageHandler = shared.MessageHandler
|
||||
|
||||
type Preset int
|
||||
|
||||
const (
|
||||
PresetAgent Preset = iota
|
||||
PresetRelay
|
||||
PresetFlared
|
||||
)
|
||||
|
||||
type presetConfig struct {
|
||||
HeaderKey string
|
||||
WSPath string
|
||||
}
|
||||
|
||||
var presets = map[Preset]presetConfig{
|
||||
PresetAgent: {HeaderKey: "X-Agent-Token", WSPath: "/api/v1/agent/ws"},
|
||||
PresetRelay: {HeaderKey: "X-Agent-Token", WSPath: "/api/v1/relay/ws"},
|
||||
PresetFlared: {HeaderKey: "X-Tunnel-Token", WSPath: "/api/v1/tunnel/ws"},
|
||||
}
|
||||
|
||||
func PresetHeaderKey(preset Preset) string {
|
||||
return presets[preset].HeaderKey
|
||||
}
|
||||
|
||||
func PresetWSPath(preset Preset) string {
|
||||
return presets[preset].WSPath
|
||||
}
|
||||
|
||||
type Client struct {
|
||||
sharedClient *shared.Client
|
||||
}
|
||||
|
||||
func New(preset Preset, baseURL, token string, timeout time.Duration) *Client {
|
||||
cfg := presets[preset]
|
||||
return &Client{
|
||||
sharedClient: shared.New(shared.Config{
|
||||
BaseURL: baseURL,
|
||||
Token: token,
|
||||
Timeout: timeout,
|
||||
HeaderKey: cfg.HeaderKey,
|
||||
WSPath: cfg.WSPath,
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Client) SetToken(token string) {
|
||||
c.sharedClient.SetToken(token)
|
||||
}
|
||||
|
||||
func (c *Client) URL() string {
|
||||
return c.sharedClient.URL()
|
||||
}
|
||||
|
||||
type Connection struct {
|
||||
sharedConn *shared.Connection
|
||||
}
|
||||
|
||||
func (c *Client) Connect(ctx context.Context) (*Connection, error) {
|
||||
conn, err := c.sharedClient.Connect(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &Connection{sharedConn: conn}, nil
|
||||
}
|
||||
|
||||
type AgentConnection struct {
|
||||
Connection
|
||||
}
|
||||
|
||||
func (c *Client) ConnectAgent(ctx context.Context) (*AgentConnection, error) {
|
||||
conn, err := c.Connect(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &AgentConnection{Connection: *conn}, nil
|
||||
}
|
||||
|
||||
func (conn *Connection) URL() string {
|
||||
if conn == nil || conn.sharedConn == nil {
|
||||
return ""
|
||||
}
|
||||
return conn.sharedConn.URL
|
||||
}
|
||||
|
||||
func (conn *Connection) SendPing() error {
|
||||
return conn.sharedConn.SendMessage(pkgprotocol.WSMessageTypePing, nil)
|
||||
}
|
||||
|
||||
func (conn *Connection) SendPong() error {
|
||||
return conn.sharedConn.SendMessage(pkgprotocol.WSMessageTypePong, nil)
|
||||
}
|
||||
|
||||
func (conn *Connection) SendMessage(msgType string, payload any) error {
|
||||
return conn.sharedConn.SendMessage(msgType, payload)
|
||||
}
|
||||
|
||||
func (conn *Connection) Receive() (pkgprotocol.WSMessage, error) {
|
||||
var message pkgprotocol.WSMessage
|
||||
if err := conn.sharedConn.Receive(&message); err != nil {
|
||||
return message, err
|
||||
}
|
||||
return message, nil
|
||||
}
|
||||
|
||||
func (conn *Connection) RunReceiveLoop(ctx context.Context, handler MessageHandler) error {
|
||||
return conn.sharedConn.RunReceiveLoop(ctx, handler)
|
||||
}
|
||||
|
||||
func (conn *Connection) Close() error {
|
||||
if conn == nil || conn.sharedConn == nil {
|
||||
return nil
|
||||
}
|
||||
return conn.sharedConn.Close()
|
||||
}
|
||||
|
||||
func (conn *AgentConnection) SendStatus(payload pkgprotocol.NodePayload) error {
|
||||
return conn.sharedConn.SendMessage(pkgprotocol.WSMessageTypeStatus, payload)
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
package wsclient
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestPresetConfig(t *testing.T) {
|
||||
tests := []struct {
|
||||
preset Preset
|
||||
headerKey string
|
||||
wsPath string
|
||||
}{
|
||||
{preset: PresetAgent, headerKey: "X-Agent-Token", wsPath: "/api/v1/agent/ws"},
|
||||
{preset: PresetRelay, headerKey: "X-Agent-Token", wsPath: "/api/v1/relay/ws"},
|
||||
{preset: PresetFlared, headerKey: "X-Tunnel-Token", wsPath: "/api/v1/tunnel/ws"},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.wsPath, func(t *testing.T) {
|
||||
if got := PresetHeaderKey(tt.preset); got != tt.headerKey {
|
||||
t.Fatalf("PresetHeaderKey() = %q, want %q", got, tt.headerKey)
|
||||
}
|
||||
if got := PresetWSPath(tt.preset); got != tt.wsPath {
|
||||
t.Fatalf("PresetWSPath() = %q, want %q", got, tt.wsPath)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestClientURL(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
preset Preset
|
||||
baseURL string
|
||||
want string
|
||||
}{
|
||||
{
|
||||
name: "agent https",
|
||||
preset: PresetAgent,
|
||||
baseURL: "https://example.com",
|
||||
want: "wss://example.com/api/v1/agent/ws",
|
||||
},
|
||||
{
|
||||
name: "relay http with path prefix",
|
||||
preset: PresetRelay,
|
||||
baseURL: "http://example.com/api",
|
||||
want: "ws://example.com/api/api/v1/relay/ws",
|
||||
},
|
||||
{
|
||||
name: "flared wss",
|
||||
preset: PresetFlared,
|
||||
baseURL: "wss://edge.example.com",
|
||||
want: "wss://edge.example.com/api/v1/tunnel/ws",
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
client := New(tt.preset, tt.baseURL, "token", time.Second)
|
||||
if got := client.URL(); got != tt.want {
|
||||
t.Fatalf("URL() = %q, want %q", got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -2,74 +2,29 @@ package wsclient
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
service "github.com/Rain-kl/Wavelet/pkg/protocol"
|
||||
shared "github.com/Rain-kl/Wavelet/pkg/wsclient"
|
||||
edgews "github.com/Rain-kl/Wavelet/internal/apps/edge/wsclient"
|
||||
)
|
||||
|
||||
type WSMessage = shared.WSMessage
|
||||
type MessageHandler = shared.MessageHandler
|
||||
type WSMessage = edgews.WSMessage
|
||||
type MessageHandler = edgews.MessageHandler
|
||||
type Connection = edgews.Connection
|
||||
|
||||
type Client struct {
|
||||
sharedClient *shared.Client
|
||||
inner *edgews.Client
|
||||
}
|
||||
|
||||
type Connection struct {
|
||||
sharedConn *shared.Connection
|
||||
}
|
||||
|
||||
func New(baseURL string, token string, timeout time.Duration) *Client {
|
||||
func New(baseURL, token string, timeout time.Duration) *Client {
|
||||
return &Client{
|
||||
sharedClient: shared.New(shared.Config{
|
||||
BaseURL: baseURL,
|
||||
Token: token,
|
||||
Timeout: timeout,
|
||||
HeaderKey: "X-Tunnel-Token",
|
||||
WSPath: "/api/v1/tunnel/ws",
|
||||
}),
|
||||
inner: edgews.New(edgews.PresetFlared, baseURL, token, timeout),
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Client) SetToken(token string) {
|
||||
c.sharedClient.SetToken(token)
|
||||
c.inner.SetToken(token)
|
||||
}
|
||||
|
||||
func (c *Client) Connect(ctx context.Context) (*Connection, error) {
|
||||
conn, err := c.sharedClient.Connect(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &Connection{sharedConn: conn}, nil
|
||||
}
|
||||
|
||||
func (conn *Connection) SendPing() error {
|
||||
return conn.sharedConn.SendMessage("ping", nil)
|
||||
}
|
||||
|
||||
func (conn *Connection) SendPong() error {
|
||||
return conn.sharedConn.SendMessage("pong", nil)
|
||||
}
|
||||
|
||||
func (conn *Connection) Receive() (service.WSMessage, error) {
|
||||
var raw struct {
|
||||
Type string `json:"type"`
|
||||
Payload json.RawMessage `json:"payload,omitempty"`
|
||||
}
|
||||
if err := conn.sharedConn.Receive(&raw); err != nil {
|
||||
return service.WSMessage{}, err
|
||||
}
|
||||
return service.WSMessage{
|
||||
Type: raw.Type,
|
||||
Payload: raw.Payload,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (conn *Connection) RunReceiveLoop(ctx context.Context, handler shared.MessageHandler) error {
|
||||
return conn.sharedConn.RunReceiveLoop(ctx, handler)
|
||||
}
|
||||
|
||||
func (conn *Connection) Close() error {
|
||||
return conn.sharedConn.Close()
|
||||
}
|
||||
return c.inner.Connect(ctx)
|
||||
}
|
||||
@@ -29,82 +29,6 @@ const (
|
||||
accessLogPathMaxLength = 100
|
||||
)
|
||||
|
||||
// NodeSystemProfile is the agent-reported system profile.
|
||||
type NodeSystemProfile struct {
|
||||
Hostname string `json:"hostname"`
|
||||
OSName string `json:"os_name"`
|
||||
OSVersion string `json:"os_version"`
|
||||
KernelVersion string `json:"kernel_version"`
|
||||
Architecture string `json:"architecture"`
|
||||
CPUModel string `json:"cpu_model"`
|
||||
CPUCores int `json:"cpu_cores"`
|
||||
TotalMemoryBytes int64 `json:"total_memory_bytes"`
|
||||
TotalDiskBytes int64 `json:"total_disk_bytes"`
|
||||
UptimeSeconds int64 `json:"uptime_seconds"`
|
||||
ReportedAtUnix int64 `json:"reported_at_unix"`
|
||||
}
|
||||
|
||||
// NodeMetricSnapshot is the agent-reported capacity snapshot.
|
||||
type NodeMetricSnapshot struct {
|
||||
CapturedAtUnix int64 `json:"captured_at_unix"`
|
||||
CPUUsagePercent float64 `json:"cpu_usage_percent"`
|
||||
MemoryUsedBytes int64 `json:"memory_used_bytes"`
|
||||
MemoryTotalBytes int64 `json:"memory_total_bytes"`
|
||||
StorageUsedBytes int64 `json:"storage_used_bytes"`
|
||||
StorageTotalBytes int64 `json:"storage_total_bytes"`
|
||||
DiskReadBytes int64 `json:"disk_read_bytes"`
|
||||
DiskWriteBytes int64 `json:"disk_write_bytes"`
|
||||
NetworkRxBytes int64 `json:"network_rx_bytes"`
|
||||
NetworkTxBytes int64 `json:"network_tx_bytes"`
|
||||
}
|
||||
|
||||
// NodeOpenrestyObservation is the agent-reported openresty network observation.
|
||||
type NodeOpenrestyObservation struct {
|
||||
CapturedAtUnix int64 `json:"captured_at_unix"`
|
||||
OpenrestyRxBytes int64 `json:"openresty_rx_bytes"`
|
||||
OpenrestyTxBytes int64 `json:"openresty_tx_bytes"`
|
||||
OpenrestyConnections int64 `json:"openresty_connections"`
|
||||
}
|
||||
|
||||
// NodeTrafficReport is the agent-reported traffic window.
|
||||
type NodeTrafficReport struct {
|
||||
WindowStartedAtUnix int64 `json:"window_started_at_unix"`
|
||||
WindowEndedAtUnix int64 `json:"window_ended_at_unix"`
|
||||
RequestCount int64 `json:"request_count"`
|
||||
ErrorCount int64 `json:"error_count"`
|
||||
UniqueVisitorCount int64 `json:"unique_visitor_count"`
|
||||
StatusCodes map[string]int64 `json:"status_codes"`
|
||||
TopDomains map[string]int64 `json:"top_domains"`
|
||||
SourceCountries map[string]int64 `json:"source_countries"`
|
||||
}
|
||||
|
||||
// NodeAccessLog is a single access log row from the agent.
|
||||
type NodeAccessLog struct {
|
||||
LoggedAtUnix int64 `json:"logged_at_unix"`
|
||||
RemoteAddr string `json:"remote_addr"`
|
||||
Host string `json:"host"`
|
||||
Path string `json:"path"`
|
||||
StatusCode int `json:"status_code"`
|
||||
}
|
||||
|
||||
// BufferedObservabilityRecord is a buffered observability window from the agent.
|
||||
type BufferedObservabilityRecord struct {
|
||||
WindowStartedAtUnix int64 `json:"window_started_at_unix"`
|
||||
Snapshot *NodeMetricSnapshot `json:"snapshot,omitempty"`
|
||||
OpenrestyObservation *NodeOpenrestyObservation `json:"openresty_observation,omitempty"`
|
||||
TrafficReport *NodeTrafficReport `json:"traffic_report,omitempty"`
|
||||
AccessLogs []NodeAccessLog `json:"access_logs,omitempty"`
|
||||
}
|
||||
|
||||
// NodeHealthEvent is an agent-reported health event.
|
||||
type NodeHealthEvent struct {
|
||||
EventType string `json:"event_type"`
|
||||
Severity string `json:"severity"`
|
||||
Message string `json:"message"`
|
||||
TriggeredAtUnix int64 `json:"triggered_at_unix"`
|
||||
Metadata map[string]string `json:"metadata"`
|
||||
}
|
||||
|
||||
// PersistHeartbeatObservability stores profile, snapshots, traffic, access logs, and health events.
|
||||
func PersistHeartbeatObservability(ctx context.Context, nodeID string, payload NodePayload, reportedAt time.Time) {
|
||||
if strings.TrimSpace(nodeID) == "" {
|
||||
|
||||
@@ -0,0 +1,26 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package agent
|
||||
|
||||
import pkgprotocol "github.com/Rain-kl/Wavelet/pkg/protocol"
|
||||
|
||||
type NodePayload = pkgprotocol.NodePayload
|
||||
type NodeSystemProfile = pkgprotocol.NodeSystemProfile
|
||||
type NodeMetricSnapshot = pkgprotocol.NodeMetricSnapshot
|
||||
type NodeOpenrestyObservation = pkgprotocol.NodeOpenrestyObservation
|
||||
type NodeTrafficReport = pkgprotocol.NodeTrafficReport
|
||||
type NodeAccessLog = pkgprotocol.NodeAccessLog
|
||||
type BufferedObservabilityRecord = pkgprotocol.BufferedObservabilityRecord
|
||||
type NodeHealthEvent = pkgprotocol.NodeHealthEvent
|
||||
type ApplyLogPayload = pkgprotocol.ApplyLogPayload
|
||||
type Settings = pkgprotocol.AgentSettings
|
||||
type ActiveConfigMeta = pkgprotocol.ActiveConfigMeta
|
||||
type SupportFile = pkgprotocol.SupportFile
|
||||
type WAFIPGroup = pkgprotocol.WAFIPGroup
|
||||
type WAFIPGroupSyncRequest = pkgprotocol.WAFIPGroupSyncRequest
|
||||
type WAFIPGroupSyncResponse = pkgprotocol.WAFIPGroupSyncResponse
|
||||
|
||||
// Backward-compatible names used by server routers and handlers.
|
||||
type WAFIPGroupSyncInput = WAFIPGroupSyncRequest
|
||||
type WAFIPGroupSyncResult = WAFIPGroupSyncResponse
|
||||
@@ -16,71 +16,16 @@ const (
|
||||
applyResultFailed = "failed"
|
||||
)
|
||||
|
||||
// NodePayload is the agent register/heartbeat payload.
|
||||
type NodePayload struct {
|
||||
NodeID string `json:"node_id"`
|
||||
Name string `json:"name"`
|
||||
IP string `json:"ip"`
|
||||
Version string `json:"version"`
|
||||
ExtVersion string `json:"ext_version"`
|
||||
CurrentVersion string `json:"current_version"`
|
||||
LastError string `json:"last_error"`
|
||||
OpenrestyStatus string `json:"openresty_status"`
|
||||
OpenrestyMessage string `json:"openresty_message"`
|
||||
Profile *NodeSystemProfile `json:"profile,omitempty"`
|
||||
Snapshot *NodeMetricSnapshot `json:"snapshot,omitempty"`
|
||||
OpenrestyObservation *NodeOpenrestyObservation `json:"openresty_observation,omitempty"`
|
||||
TrafficReport *NodeTrafficReport `json:"traffic_report,omitempty"`
|
||||
AccessLogs []NodeAccessLog `json:"access_logs,omitempty"`
|
||||
BufferedObservability []BufferedObservabilityRecord `json:"buffered_observability,omitempty"`
|
||||
HealthEvents []NodeHealthEvent `json:"health_events"`
|
||||
WAFIPGroupChecksums map[string]string `json:"waf_ip_group_checksums,omitempty"`
|
||||
}
|
||||
|
||||
// ApplyLogPayload is the agent apply log report payload.
|
||||
type ApplyLogPayload struct {
|
||||
NodeID string `json:"node_id"`
|
||||
Version string `json:"version"`
|
||||
Result string `json:"result"`
|
||||
Message string `json:"message"`
|
||||
Checksum string `json:"checksum"`
|
||||
MainConfigChecksum string `json:"main_config_checksum"`
|
||||
RouteConfigChecksum string `json:"route_config_checksum"`
|
||||
SupportFileCount int `json:"support_file_count"`
|
||||
}
|
||||
|
||||
// RegistrationResponse is returned after agent registration.
|
||||
// Server uses access_token; the agent client expects agent_token via RegisterNodeResponse.
|
||||
type RegistrationResponse struct {
|
||||
NodeID string `json:"node_id"`
|
||||
AccessToken string `json:"access_token"`
|
||||
Name string `json:"name"`
|
||||
}
|
||||
|
||||
// Settings carries remote agent control flags.
|
||||
type Settings struct {
|
||||
HeartbeatInterval int `json:"heartbeat_interval"`
|
||||
WebsocketUpgradeEnabled bool `json:"websocket_upgrade_enabled"`
|
||||
AutoUpdate bool `json:"auto_update"`
|
||||
UpdateRepo string `json:"update_repo"`
|
||||
UpdateNow bool `json:"update_now"`
|
||||
UpdateChannel string `json:"update_channel"`
|
||||
UpdateTag string `json:"update_tag"`
|
||||
RestartOpenrestyNow bool `json:"restart_openresty_now"`
|
||||
}
|
||||
|
||||
// ActiveConfigMeta summarizes the active configuration version.
|
||||
type ActiveConfigMeta struct {
|
||||
Version string `json:"version"`
|
||||
Checksum string `json:"checksum"`
|
||||
}
|
||||
|
||||
// SupportFile is a configuration support artifact shipped to agents.
|
||||
type SupportFile struct {
|
||||
Path string `json:"path"`
|
||||
Content string `json:"content"`
|
||||
}
|
||||
|
||||
// ConfigResponse is the full active config payload for agents.
|
||||
// Server uses time.Time for CreatedAt; the agent client uses string via ActiveConfigResponse.
|
||||
type ConfigResponse struct {
|
||||
Version string `json:"version"`
|
||||
Checksum string `json:"checksum"`
|
||||
@@ -89,31 +34,10 @@ type ConfigResponse struct {
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
}
|
||||
|
||||
// WAFIPGroup is a WAF IP group snapshot for agents.
|
||||
type WAFIPGroup struct {
|
||||
ID uint `json:"id"`
|
||||
Name string `json:"name"`
|
||||
Type string `json:"type"`
|
||||
Enabled bool `json:"enabled"`
|
||||
IPList []string `json:"ip_list"`
|
||||
Checksum string `json:"checksum"`
|
||||
}
|
||||
|
||||
// WAFIPGroupSyncInput requests changed WAF IP groups.
|
||||
type WAFIPGroupSyncInput struct {
|
||||
IDs []uint `json:"ids"`
|
||||
Checksums map[string]string `json:"checksums"`
|
||||
}
|
||||
|
||||
// WAFIPGroupSyncResult returns synced WAF IP groups.
|
||||
type WAFIPGroupSyncResult struct {
|
||||
Groups []WAFIPGroup `json:"groups"`
|
||||
}
|
||||
|
||||
// HeartbeatResponse is the heartbeat handler result.
|
||||
type HeartbeatResponse struct {
|
||||
Node *model.OpenFlareNode `json:"node"`
|
||||
AgentSettings *Settings `json:"agent_settings"`
|
||||
ActiveConfig *ActiveConfigMeta `json:"active_config"`
|
||||
WAFIPGroups []WAFIPGroup `json:"waf_ip_groups,omitempty"`
|
||||
}
|
||||
}
|
||||
@@ -11,7 +11,6 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/openflare/agent"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/openflare/relay"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"gorm.io/gorm"
|
||||
@@ -24,73 +23,6 @@ const (
|
||||
applyResultFail = "failed"
|
||||
)
|
||||
|
||||
// HeartbeatPayload is sent by OpenFlared on each heartbeat.
|
||||
type HeartbeatPayload struct {
|
||||
ClientVersion string `json:"client_version"`
|
||||
FrpVersion string `json:"frp_version"`
|
||||
IP string `json:"ip"`
|
||||
TunnelStatus string `json:"tunnel_status"`
|
||||
ConnectedRelays []ConnectedRelay `json:"connected_relays"`
|
||||
CurrentVersion string `json:"current_version"`
|
||||
CurrentChecksum string `json:"current_checksum"`
|
||||
}
|
||||
|
||||
// ConnectedRelay describes relay connection status from the client.
|
||||
type ConnectedRelay struct {
|
||||
RelayNodeID string `json:"relay_node_id"`
|
||||
Status string `json:"status"`
|
||||
ProxyCount int `json:"proxy_count"`
|
||||
}
|
||||
|
||||
// ActiveConfigMeta summarizes the active config version.
|
||||
type ActiveConfigMeta struct {
|
||||
Version string `json:"version"`
|
||||
Checksum string `json:"checksum"`
|
||||
}
|
||||
|
||||
// HeartbeatResponse is returned to the OpenFlared client.
|
||||
type HeartbeatResponse struct {
|
||||
ActiveConfig *ActiveConfigMeta `json:"active_config"`
|
||||
TunnelSettings *relay.Settings `json:"tunnel_settings"`
|
||||
}
|
||||
|
||||
// TunnelConfigResponse is the full tunnel routing config sent to the client.
|
||||
type TunnelConfigResponse struct {
|
||||
Version string `json:"version"`
|
||||
Checksum string `json:"checksum"`
|
||||
Relays []RelayInfo `json:"relays"`
|
||||
Proxies []ProxyEntry `json:"proxies"`
|
||||
}
|
||||
|
||||
// RelayInfo describes a relay the client should connect to.
|
||||
type RelayInfo struct {
|
||||
RelayNodeID string `json:"relay_node_id"`
|
||||
Address string `json:"address"`
|
||||
AuthToken string `json:"auth_token"`
|
||||
ProxyURL string `json:"proxy_url"`
|
||||
}
|
||||
|
||||
// ProxyEntry describes one frpc proxy definition.
|
||||
type ProxyEntry struct {
|
||||
Name string `json:"name"`
|
||||
Type string `json:"type"`
|
||||
LocalAddr string `json:"local_addr"`
|
||||
LocalPort int `json:"local_port"`
|
||||
CustomDomains []string `json:"custom_domains"`
|
||||
}
|
||||
|
||||
// ApplyLogPayload is the apply result reported by OpenFlared.
|
||||
type ApplyLogPayload struct {
|
||||
NodeID string `json:"node_id"`
|
||||
Version string `json:"version"`
|
||||
Result string `json:"result"`
|
||||
Message string `json:"message"`
|
||||
Checksum string `json:"checksum"`
|
||||
MainConfigChecksum string `json:"main_config_checksum"`
|
||||
RouteConfigChecksum string `json:"route_config_checksum"`
|
||||
SupportFileCount int `json:"support_file_count"`
|
||||
}
|
||||
|
||||
// Heartbeat processes an OpenFlared heartbeat and returns runtime settings.
|
||||
func Heartbeat(ctx context.Context, node *model.OpenFlareNode, payload HeartbeatPayload) (*HeartbeatResponse, error) {
|
||||
if node == nil {
|
||||
|
||||
@@ -0,0 +1,15 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package flared
|
||||
|
||||
import pkgprotocol "github.com/Rain-kl/Wavelet/pkg/protocol"
|
||||
|
||||
type HeartbeatPayload = pkgprotocol.FlaredHeartbeatPayload
|
||||
type ConnectedRelay = pkgprotocol.FlaredConnectedRelay
|
||||
type ActiveConfigMeta = pkgprotocol.ActiveConfigMeta
|
||||
type HeartbeatResponse = pkgprotocol.FlaredHeartbeatResponse
|
||||
type TunnelConfigResponse = pkgprotocol.FlaredTunnelConfigResponse
|
||||
type RelayInfo = pkgprotocol.FlaredRelayInfo
|
||||
type ProxyEntry = pkgprotocol.FlaredProxyEntry
|
||||
type ApplyLogPayload = pkgprotocol.ApplyLogPayload
|
||||
@@ -16,59 +16,6 @@ import (
|
||||
|
||||
const nodeStatusOnline = "online"
|
||||
|
||||
// ProxyStat describes a single frps proxy reported by the relay.
|
||||
type ProxyStat struct {
|
||||
Name string `json:"name"`
|
||||
Type string `json:"type"`
|
||||
Status string `json:"status"`
|
||||
ClientVersion string `json:"client_version"`
|
||||
LastStartTime string `json:"last_start_time"`
|
||||
LastCloseTime string `json:"last_close_time"`
|
||||
ClientAddr string `json:"client_addr"`
|
||||
}
|
||||
|
||||
// HeartbeatPayload is sent by OpenFlareRelay on each heartbeat.
|
||||
type HeartbeatPayload struct {
|
||||
Version string `json:"version"`
|
||||
ExtVersion string `json:"frp_version"`
|
||||
RelayStatus string `json:"relay_status"`
|
||||
FrpsConnCount int `json:"frps_connections"`
|
||||
FrpsProxyCount int `json:"frps_proxy_count"`
|
||||
FrpsClientCount int `json:"frps_client_count"`
|
||||
FrpsProxies []ProxyStat `json:"frps_proxies,omitempty"`
|
||||
Name string `json:"name"`
|
||||
IP string `json:"ip"`
|
||||
Profile *agent.NodeSystemProfile `json:"profile,omitempty"`
|
||||
Snapshot *agent.NodeMetricSnapshot `json:"snapshot,omitempty"`
|
||||
HealthEvents []agent.NodeHealthEvent `json:"health_events,omitempty"`
|
||||
}
|
||||
|
||||
// Config is the frps configuration sent to the relay.
|
||||
type Config struct {
|
||||
BindPort int `json:"bind_port"`
|
||||
VhostHTTPPort int `json:"vhost_http_port"`
|
||||
AuthToken string `json:"auth_token"`
|
||||
LogLevel string `json:"log_level"`
|
||||
WebServerEnabled bool `json:"web_server_enabled"`
|
||||
}
|
||||
|
||||
// Settings contains runtime settings for relay and flared clients.
|
||||
type Settings struct {
|
||||
HeartbeatInterval int `json:"heartbeat_interval"`
|
||||
WebsocketUpgradeEnabled bool `json:"websocket_upgrade_enabled"`
|
||||
AutoUpdate bool `json:"auto_update"`
|
||||
UpdateRepo string `json:"update_repo"`
|
||||
UpdateNow bool `json:"update_now"`
|
||||
UpdateChannel string `json:"update_channel"`
|
||||
UpdateTag string `json:"update_tag"`
|
||||
}
|
||||
|
||||
// HeartbeatResponse is returned from a relay heartbeat.
|
||||
type HeartbeatResponse struct {
|
||||
RelayConfig *Config `json:"relay_config"`
|
||||
RelaySettings *Settings `json:"relay_settings"`
|
||||
}
|
||||
|
||||
// Heartbeat processes a relay heartbeat, updates node status, and returns config.
|
||||
func Heartbeat(ctx context.Context, node *model.OpenFlareNode, payload HeartbeatPayload) (*HeartbeatResponse, error) {
|
||||
if node == nil {
|
||||
|
||||
@@ -0,0 +1,12 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package relay
|
||||
|
||||
import pkgprotocol "github.com/Rain-kl/Wavelet/pkg/protocol"
|
||||
|
||||
type ProxyStat = pkgprotocol.RelayProxyStat
|
||||
type HeartbeatPayload = pkgprotocol.RelayHeartbeatPayload
|
||||
type Config = pkgprotocol.RelayConfig
|
||||
type Settings = pkgprotocol.RelaySettings
|
||||
type HeartbeatResponse = pkgprotocol.RelayHeartbeatResponse
|
||||
@@ -2,74 +2,29 @@ package wsclient
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
service "github.com/Rain-kl/Wavelet/pkg/protocol"
|
||||
shared "github.com/Rain-kl/Wavelet/pkg/wsclient"
|
||||
edgews "github.com/Rain-kl/Wavelet/internal/apps/edge/wsclient"
|
||||
)
|
||||
|
||||
type WSMessage = shared.WSMessage
|
||||
type MessageHandler = shared.MessageHandler
|
||||
type WSMessage = edgews.WSMessage
|
||||
type MessageHandler = edgews.MessageHandler
|
||||
type Connection = edgews.Connection
|
||||
|
||||
type Client struct {
|
||||
sharedClient *shared.Client
|
||||
inner *edgews.Client
|
||||
}
|
||||
|
||||
type Connection struct {
|
||||
sharedConn *shared.Connection
|
||||
}
|
||||
|
||||
func New(baseURL string, token string, timeout time.Duration) *Client {
|
||||
func New(baseURL, token string, timeout time.Duration) *Client {
|
||||
return &Client{
|
||||
sharedClient: shared.New(shared.Config{
|
||||
BaseURL: baseURL,
|
||||
Token: token,
|
||||
Timeout: timeout,
|
||||
HeaderKey: "X-Agent-Token",
|
||||
WSPath: "/api/v1/relay/ws",
|
||||
}),
|
||||
inner: edgews.New(edgews.PresetRelay, baseURL, token, timeout),
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Client) SetToken(token string) {
|
||||
c.sharedClient.SetToken(token)
|
||||
c.inner.SetToken(token)
|
||||
}
|
||||
|
||||
func (c *Client) Connect(ctx context.Context) (*Connection, error) {
|
||||
conn, err := c.sharedClient.Connect(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &Connection{sharedConn: conn}, nil
|
||||
}
|
||||
|
||||
func (conn *Connection) SendPing() error {
|
||||
return conn.sharedConn.SendMessage("ping", nil)
|
||||
}
|
||||
|
||||
func (conn *Connection) SendPong() error {
|
||||
return conn.sharedConn.SendMessage("pong", nil)
|
||||
}
|
||||
|
||||
func (conn *Connection) Receive() (service.WSMessage, error) {
|
||||
var raw struct {
|
||||
Type string `json:"type"`
|
||||
Payload json.RawMessage `json:"payload,omitempty"`
|
||||
}
|
||||
if err := conn.sharedConn.Receive(&raw); err != nil {
|
||||
return service.WSMessage{}, err
|
||||
}
|
||||
return service.WSMessage{
|
||||
Type: raw.Type,
|
||||
Payload: raw.Payload,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (conn *Connection) RunReceiveLoop(ctx context.Context, handler shared.MessageHandler) error {
|
||||
return conn.sharedConn.RunReceiveLoop(ctx, handler)
|
||||
}
|
||||
|
||||
func (conn *Connection) Close() error {
|
||||
return conn.sharedConn.Close()
|
||||
}
|
||||
return c.inner.Connect(ctx)
|
||||
}
|
||||
+3
-34
@@ -1,39 +1,8 @@
|
||||
package protocol
|
||||
|
||||
type AgentNodeSystemProfile struct {
|
||||
Hostname string `json:"hostname"`
|
||||
OSName string `json:"os_name"`
|
||||
OSVersion string `json:"os_version"`
|
||||
KernelVersion string `json:"kernel_version"`
|
||||
Architecture string `json:"architecture"`
|
||||
CPUModel string `json:"cpu_model"`
|
||||
CPUCores int `json:"cpu_cores"`
|
||||
TotalMemoryBytes int64 `json:"total_memory_bytes"`
|
||||
TotalDiskBytes int64 `json:"total_disk_bytes"`
|
||||
UptimeSeconds int64 `json:"uptime_seconds"`
|
||||
ReportedAtUnix int64 `json:"reported_at_unix"`
|
||||
}
|
||||
|
||||
type AgentNodeMetricSnapshot struct {
|
||||
CapturedAtUnix int64 `json:"captured_at_unix"`
|
||||
CPUUsagePercent float64 `json:"cpu_usage_percent"`
|
||||
MemoryUsedBytes int64 `json:"memory_used_bytes"`
|
||||
MemoryTotalBytes int64 `json:"memory_total_bytes"`
|
||||
StorageUsedBytes int64 `json:"storage_used_bytes"`
|
||||
StorageTotalBytes int64 `json:"storage_total_bytes"`
|
||||
DiskReadBytes int64 `json:"disk_read_bytes"`
|
||||
DiskWriteBytes int64 `json:"disk_write_bytes"`
|
||||
NetworkRxBytes int64 `json:"network_rx_bytes"`
|
||||
NetworkTxBytes int64 `json:"network_tx_bytes"`
|
||||
}
|
||||
|
||||
type AgentNodeHealthEvent struct {
|
||||
EventType string `json:"event_type"`
|
||||
Severity string `json:"severity"`
|
||||
Message string `json:"message"`
|
||||
TriggeredAtUnix int64 `json:"triggered_at_unix"`
|
||||
Metadata map[string]string `json:"metadata"`
|
||||
}
|
||||
type AgentNodeSystemProfile = NodeSystemProfile
|
||||
type AgentNodeMetricSnapshot = NodeMetricSnapshot
|
||||
type AgentNodeHealthEvent = NodeHealthEvent
|
||||
|
||||
type RelayProxyStat struct {
|
||||
Name string `json:"name"`
|
||||
|
||||
Reference in New Issue
Block a user