From 0b3479270944735f0fae8a5726503ba862d746de Mon Sep 17 00:00:00 2001 From: ryan Date: Fri, 19 Jun 2026 15:00:46 +0800 Subject: [PATCH] =?UTF-8?q?refactor(edge):=20Server=20=E5=8D=8F=E8=AE=AE?= =?UTF-8?q?=E7=BB=9F=E4=B8=80=E4=B8=8E=20wsclient=20=E6=94=B6=E6=95=9B?= =?UTF-8?q?=EF=BC=88Phase=203=20Batch=203=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - openflare/agent/relay/flared 协议类型改为 pkg/protocol 别名,删除重复 struct - 新增 edge/wsclient Preset 表,三组件 wsclient 改为薄包装 - 更新设计文档、计划与 changelog --- docs/changelog/index.md | 2 + docs/design/edge-runtime-refactor.md | 18 ++- docs/plan/20260619-edge-phase3-tasks.md | 28 +++- internal/apps/agent/wsclient/client.go | 69 ++-------- internal/apps/edge/wsclient/client.go | 130 ++++++++++++++++++ internal/apps/edge/wsclient/client_test.go | 66 +++++++++ internal/apps/flared/wsclient/client.go | 65 ++------- .../apps/openflare/agent/observability.go | 76 ---------- .../apps/openflare/agent/protocol_alias.go | 26 ++++ internal/apps/openflare/agent/types.go | 82 +---------- internal/apps/openflare/flared/logics.go | 68 --------- .../apps/openflare/flared/protocol_alias.go | 15 ++ internal/apps/openflare/relay/logics.go | 53 ------- .../apps/openflare/relay/protocol_alias.go | 12 ++ internal/apps/relay/wsclient/client.go | 65 ++------- pkg/protocol/server.go | 37 +---- 16 files changed, 328 insertions(+), 484 deletions(-) create mode 100644 internal/apps/edge/wsclient/client.go create mode 100644 internal/apps/edge/wsclient/client_test.go create mode 100644 internal/apps/openflare/agent/protocol_alias.go create mode 100644 internal/apps/openflare/flared/protocol_alias.go create mode 100644 internal/apps/openflare/relay/protocol_alias.go diff --git a/docs/changelog/index.md b/docs/changelog/index.md index 234ba25d..51928a18 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -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: diff --git a/docs/design/edge-runtime-refactor.md b/docs/design/edge-runtime-refactor.md index 0e00a693..0b5708ac 100644 --- a/docs/design/edge-runtime-refactor.md +++ b/docs/design/edge-runtime-refactor.md @@ -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 配置表化 | --- diff --git a/docs/plan/20260619-edge-phase3-tasks.md b/docs/plan/20260619-edge-phase3-tasks.md index 29554905..175daf05 100644 --- a/docs/plan/20260619-edge-phase3-tasks.md +++ b/docs/plan/20260619-edge-phase3-tasks.md @@ -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 ``` \ No newline at end of file diff --git a/internal/apps/agent/wsclient/client.go b/internal/apps/agent/wsclient/client.go index c71750f1..5440ed74 100644 --- a/internal/apps/agent/wsclient/client.go +++ b/internal/apps/agent/wsclient/client.go @@ -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) +} \ No newline at end of file diff --git a/internal/apps/edge/wsclient/client.go b/internal/apps/edge/wsclient/client.go new file mode 100644 index 00000000..fd168e1c --- /dev/null +++ b/internal/apps/edge/wsclient/client.go @@ -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) +} \ No newline at end of file diff --git a/internal/apps/edge/wsclient/client_test.go b/internal/apps/edge/wsclient/client_test.go new file mode 100644 index 00000000..89828544 --- /dev/null +++ b/internal/apps/edge/wsclient/client_test.go @@ -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) + } + }) + } +} \ No newline at end of file diff --git a/internal/apps/flared/wsclient/client.go b/internal/apps/flared/wsclient/client.go index 66855de3..b4b9ddc8 100644 --- a/internal/apps/flared/wsclient/client.go +++ b/internal/apps/flared/wsclient/client.go @@ -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) +} \ No newline at end of file diff --git a/internal/apps/openflare/agent/observability.go b/internal/apps/openflare/agent/observability.go index d8acf8fc..6e7f21f7 100644 --- a/internal/apps/openflare/agent/observability.go +++ b/internal/apps/openflare/agent/observability.go @@ -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) == "" { diff --git a/internal/apps/openflare/agent/protocol_alias.go b/internal/apps/openflare/agent/protocol_alias.go new file mode 100644 index 00000000..0a2df55a --- /dev/null +++ b/internal/apps/openflare/agent/protocol_alias.go @@ -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 \ No newline at end of file diff --git a/internal/apps/openflare/agent/types.go b/internal/apps/openflare/agent/types.go index 47d6e24a..c0cf1aea 100644 --- a/internal/apps/openflare/agent/types.go +++ b/internal/apps/openflare/agent/types.go @@ -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"` -} +} \ No newline at end of file diff --git a/internal/apps/openflare/flared/logics.go b/internal/apps/openflare/flared/logics.go index 83b328c2..045a439f 100644 --- a/internal/apps/openflare/flared/logics.go +++ b/internal/apps/openflare/flared/logics.go @@ -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 { diff --git a/internal/apps/openflare/flared/protocol_alias.go b/internal/apps/openflare/flared/protocol_alias.go new file mode 100644 index 00000000..4cc0d500 --- /dev/null +++ b/internal/apps/openflare/flared/protocol_alias.go @@ -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 \ No newline at end of file diff --git a/internal/apps/openflare/relay/logics.go b/internal/apps/openflare/relay/logics.go index 79c5e3dd..2935c55c 100644 --- a/internal/apps/openflare/relay/logics.go +++ b/internal/apps/openflare/relay/logics.go @@ -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 { diff --git a/internal/apps/openflare/relay/protocol_alias.go b/internal/apps/openflare/relay/protocol_alias.go new file mode 100644 index 00000000..c130cd3a --- /dev/null +++ b/internal/apps/openflare/relay/protocol_alias.go @@ -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 \ No newline at end of file diff --git a/internal/apps/relay/wsclient/client.go b/internal/apps/relay/wsclient/client.go index db795417..e05d7c73 100644 --- a/internal/apps/relay/wsclient/client.go +++ b/internal/apps/relay/wsclient/client.go @@ -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) +} \ No newline at end of file diff --git a/pkg/protocol/server.go b/pkg/protocol/server.go index 0a12273c..8a04ca86 100644 --- a/pkg/protocol/server.go +++ b/pkg/protocol/server.go @@ -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"`