feat(relay): support configurable frps webui port and fix node detail integration

- Implement configurable FRPS WebUI switch and custom port setting (relay_frps_web_ui_port) in system configs.
- Integrate settings into Relay Node detail manage page instead of global settings.
- Dynamically query server version to select matching Docker image tag for Relay installation.
- Clean up legacy code and fix backend linter/test warnings.
This commit is contained in:
ryan
2026-06-22 12:48:54 +08:00
parent ffe98f6307
commit 346024f346
24 changed files with 204 additions and 69 deletions
@@ -27,7 +27,7 @@ import (
"github.com/Rain-kl/Wavelet/internal/common/response"
)
const expectedDefaultConfigsCount = 30
const expectedDefaultConfigsCount = 32
func setupTestRouter(authUser *model.User) *gin.Engine {
r := testhelper.NewTestGinEngine()
@@ -168,8 +168,8 @@ func TestListSystemConfigs(t *testing.T) {
var configs []model.SystemConfig
_ = json.Unmarshal(dataBytes, &configs)
if len(configs) != 1 || configs[0].Key != model.ConfigKeyMaxAPIKeysPerUser {
t.Errorf("expected 1 business config (max_api_keys_per_user), got %d: %v", len(configs), configs)
if len(configs) != 3 {
t.Errorf("expected 3 business configs, got %d: %v", len(configs), configs)
}
})
}
+23 -28
View File
@@ -68,19 +68,8 @@ func (s *Service) syncMatchingChecksum(ctx context.Context, mode string, startup
}
func (s *Service) finishUpToDateSync(ctx context.Context, mode string, snapshot *state.Snapshot, target *protocol.ActiveConfigMeta) error {
if pagesReconcileNeeded(snapshot) {
if pagesDiscoveryNeeded(snapshot) {
config, err := s.client.GetActiveConfig(ctx)
if err != nil {
slog.Error("fetch active config failed", "mode", mode, "error", err)
return err
}
if err := s.syncPagesDeployments(ctx, snapshot, config); err != nil {
return err
}
} else if err := s.syncPagesDeployments(ctx, snapshot, nil); err != nil {
return err
}
if err := s.reconcilePages(ctx, mode, snapshot); err != nil {
return err
}
slog.Debug("local openresty config already up to date", "mode", mode, "version", target.Version)
if shouldReportNoopApply(snapshot, target.Version, target.Checksum) {
@@ -96,6 +85,21 @@ func (s *Service) finishUpToDateSync(ctx context.Context, mode string, snapshot
return s.stateStore.Save(snapshot)
}
func (s *Service) reconcilePages(ctx context.Context, mode string, snapshot *state.Snapshot) error {
if !pagesReconcileNeeded(snapshot) {
return nil
}
if pagesDiscoveryNeeded(snapshot) {
config, err := s.client.GetActiveConfig(ctx)
if err != nil {
slog.Error("fetch active config failed", "mode", mode, "error", err)
return err
}
return s.syncPagesDeployments(ctx, snapshot, config)
}
return s.syncPagesDeployments(ctx, snapshot, nil)
}
func (s *Service) syncMismatchedChecksum(ctx context.Context, mode string, startup bool, snapshot *state.Snapshot, currentChecksum string, target *protocol.ActiveConfigMeta) error {
if isBlockedTarget(snapshot, target.Version, target.Checksum) {
slog.Warn("skipping blocked config version after previous failed apply", "mode", mode, "version", target.Version, "checksum", target.Checksum)
@@ -111,22 +115,13 @@ func (s *Service) syncMismatchedChecksum(ctx context.Context, mode string, start
clearBlockedTarget(snapshot)
}
if snapshot.CurrentVersion == target.Version && snapshot.CurrentChecksum == target.Checksum && !startup {
if pagesReconcileNeeded(snapshot) {
if pagesDiscoveryNeeded(snapshot) {
config, fetchErr := s.client.GetActiveConfig(ctx)
if fetchErr != nil {
slog.Error("fetch active config failed", "mode", mode, "error", fetchErr)
return fetchErr
}
if err := s.syncPagesDeployments(ctx, snapshot, config); err != nil {
return err
}
} else if err := s.syncPagesDeployments(ctx, snapshot, nil); err != nil {
return err
}
return s.stateStore.Save(snapshot)
reconciled := pagesReconcileNeeded(snapshot)
if err := s.reconcilePages(ctx, mode, snapshot); err != nil {
return err
}
if !reconciled {
slog.Debug("skipping config fetch because state already records target version/checksum", "version", target.Version, "checksum", target.Checksum)
}
slog.Debug("skipping config fetch because state already records target version/checksum", "version", target.Version, "checksum", target.Checksum)
return s.stateStore.Save(snapshot)
}
-8
View File
@@ -105,14 +105,6 @@ func applyNodeRuntime(ctx context.Context, node *model.OpenFlareNode, payload No
}
}
func cloneCoordinate(value *float64) *float64 {
if value == nil {
return nil
}
cloned := *value
return &cloned
}
func truncateForDatabase(value string, maxVal int) string {
if maxVal <= 0 {
return ""
+6 -6
View File
@@ -37,7 +37,7 @@ func EnsureRuntimeProvider(ctx context.Context) error {
runtimeInitErr = err
return
}
runtimeInitErr = applyProviderFromModel()
runtimeInitErr = applyProviderFromModel(ctx)
})
return runtimeInitErr
}
@@ -47,18 +47,18 @@ func RefreshRuntimeProvider(ctx context.Context) error {
if err := model.InitOptionMap(ctx); err != nil {
return err
}
return applyProviderFromModel()
return applyProviderFromModel(ctx)
}
func applyProviderFromModel() error {
func applyProviderFromModel(ctx context.Context) error {
model.OptionMapRWMutex.RLock()
provider := strings.TrimSpace(model.GeoIPProvider)
model.OptionMapRWMutex.RUnlock()
return ApplyProvider(provider)
return ApplyProvider(ctx, provider)
}
// ApplyProvider switches the process-wide GeoIP backend.
func ApplyProvider(provider string) error {
func ApplyProvider(ctx context.Context, provider string) error {
normalized := strings.TrimSpace(strings.ToLower(provider))
if normalized == "" {
normalized = pkggeoip.ProviderDisabled
@@ -75,7 +75,7 @@ func ApplyProvider(provider string) error {
if normalized == pkggeoip.ProviderMaxMind {
path, err := ensureServerMMDB()
if err != nil {
logger.WarnF(context.Background(), "[GeoIP] seed MaxMind database failed: %v", err)
logger.WarnF(ctx, "[GeoIP] seed MaxMind database failed: %v", err)
}
if path != "" {
pkggeoip.GeoIPFilePath = path
@@ -440,7 +440,7 @@ func matchOpenrestyObservation(
observations []*model.OpenFlareNodeObservationOpenresty,
) *model.OpenFlareNodeObservationOpenresty {
var matched *model.OpenFlareNodeObservationOpenresty
var bestDelta time.Duration = metricSnapshotOpenrestyMatchWindow + time.Second
bestDelta := metricSnapshotOpenrestyMatchWindow + time.Second
for _, observation := range observations {
if observation == nil {
continue
+8 -1
View File
@@ -4,10 +4,12 @@
package relay
import (
"context"
"net"
"strings"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
)
const (
@@ -83,16 +85,21 @@ func isPublicNodeIP(raw string) bool {
return true
}
func buildRelayConfig(node *model.OpenFlareNode) *Config {
func buildRelayConfig(ctx context.Context, node *model.OpenFlareNode) *Config {
if node == nil {
return nil
}
webServerPort, err := repository.GetIntByKey(ctx, model.ConfigKeyRelayFRPSWebUIPort)
if err != nil || webServerPort <= 0 {
webServerPort = node.RelayBindPort + 500
}
return &Config{
BindPort: node.RelayBindPort,
VhostHTTPPort: node.RelayVhostHTTPPort,
AuthToken: node.RelayAuthToken,
LogLevel: "info",
WebServerEnabled: node.RelayWebServerEnabled,
WebServerPort: webServerPort,
}
}
+1 -1
View File
@@ -98,7 +98,7 @@ func Heartbeat(ctx context.Context, node *model.OpenFlareNode, payload Heartbeat
persistRelayHeartbeatObservability(ctx, node.NodeID, payload, now)
return &HeartbeatResponse{
RelayConfig: buildRelayConfig(node),
RelayConfig: buildRelayConfig(ctx, node),
RelaySettings: BuildSettings(node, updateNow, updateChannel, updateTag),
}, nil
}
@@ -12,6 +12,7 @@ import (
"github.com/Rain-kl/Wavelet/internal/config"
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/task"
"github.com/glebarez/sqlite"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -21,11 +22,13 @@ import (
func setupSSLRenewTestDB(t *testing.T) func() {
t.Helper()
task.RegisterTaskMeta(tls.SSLSingleRenewMeta)
sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
DisableForeignKeyConstraintWhenMigrating: true,
})
require.NoError(t, err)
require.NoError(t, sqliteDB.AutoMigrate(&model.TLSCertificate{}))
require.NoError(t, sqliteDB.AutoMigrate(&model.TLSCertificate{}, &model.TaskExecution{}))
db.SetDB(sqliteDB)
oldSecret := config.Config.App.SessionSecret
+7 -2
View File
@@ -127,7 +127,8 @@ func (m *Manager) UpdateConfig(ctx context.Context, cfg *service.RelayConfig) {
m.activeConfig.BindPort == cfg.BindPort &&
m.activeConfig.VhostHTTPPort == cfg.VhostHTTPPort &&
m.activeConfig.AuthToken == cfg.AuthToken &&
m.activeConfig.WebServerEnabled == cfg.WebServerEnabled {
m.activeConfig.WebServerEnabled == cfg.WebServerEnabled &&
m.activeConfig.WebServerPort == cfg.WebServerPort {
if m.cmd == nil && !m.stopping {
slog.Warn("frps config unchanged but process is not running, restarting")
m.stopping = false
@@ -188,7 +189,11 @@ func (m *Manager) renderConfig(cfg *service.RelayConfig) error {
} else {
buf.WriteString("addr = \"127.0.0.1\"\n")
}
fmt.Fprintf(&buf, "port = %d\n", defaultFrpsWebServerPort)
port := cfg.WebServerPort
if port <= 0 {
port = defaultFrpsWebServerPort
}
fmt.Fprintf(&buf, "port = %d\n", port)
buf.WriteString("user = \"admin\"\n")
password := m.agentToken
+6
View File
@@ -53,6 +53,9 @@ exit "${EXIT_CODE:-0}"
// Helper to poll for status to eliminate timing flakiness in tests
func assertStatusEventually(t *testing.T, m *Manager, expectedStatus string, timeout time.Duration) {
if timeout < 6*time.Second {
timeout = 6 * time.Second
}
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
rt := m.GetRuntimeStatus()
@@ -67,6 +70,9 @@ func assertStatusEventually(t *testing.T, m *Manager, expectedStatus string, tim
func assertCommandExitedEventually(t *testing.T, cmd *exec.Cmd, timeout time.Duration) {
t.Helper()
if timeout < 6*time.Second {
timeout = 6 * time.Second
}
done := make(chan error, 1)
go func() {
+1 -1
View File
@@ -39,7 +39,7 @@ var (
OpenStoredUpload = ingest.OpenActiveObject
ActiveUploadHash = ingest.ActiveHash
ResolveLocalFile = ingest.ResolveLocalFile
IngestFromLocalPath = ingest.IngestFromLocalPath
IngestFromLocalPath = ingest.FromLocalPath
)
type (
+7 -3
View File
@@ -31,8 +31,8 @@ func ResolveLocalFile(ctx context.Context, req LocalFileCandidateRequest) (strin
return "", 0, os.ErrNotExist
}
// IngestFromLocalPath ingests a local regular file through the standard upload ingest path.
func IngestFromLocalPath(ctx context.Context, localPath string, req Request) (Result, error) {
// FromLocalPath ingests a local regular file through the standard upload ingest path.
func FromLocalPath(ctx context.Context, localPath string, req Request) (Result, error) {
localPath = strings.TrimSpace(localPath)
if localPath == "" {
return Result{}, errors.New("local path is required")
@@ -59,7 +59,11 @@ func IngestFromLocalPath(ctx context.Context, localPath string, req Request) (Re
func buildLocalFileCandidates(ctx context.Context, req LocalFileCandidateRequest) []string {
seen := make(map[string]struct{})
candidates := make([]string, 0, 8+len(req.RelativePaths)*4)
const (
initialCandidatesCap = 8
relativePathsWeight = 4
)
candidates := make([]string, 0, initialCandidatesCap+len(req.RelativePaths)*relativePathsWeight)
add := func(raw string) {
value := strings.TrimSpace(raw)
if value == "" {
@@ -28,8 +28,8 @@ func TestResolveLocalFileFindsStoredPath(t *testing.T) {
}
}
func TestIngestFromLocalPathRequiresPath(t *testing.T) {
_, err := IngestFromLocalPath(context.Background(), "", Request{})
func TestFromLocalPathRequiresPath(t *testing.T) {
_, err := FromLocalPath(context.Background(), "", Request{})
if err == nil {
t.Fatal("expected error for empty local path")
}