mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-04 15:06:37 +08:00
aa4faddade
Result: {"status":"keep","total_issues":8,"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_exhaustive":0,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_usestdlibvars":0,"golint_wastedassign":0,"golint_total":8,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_usetesting":0,"golint_test_total":0,"golint_vetx_total":0,"eslint_problems":0,"eslint_errors":0,"eslint_warnings":0,"tsc_errors":0,"vitest_failed":0,"vitest_total":116,"measure_s":81}
81 lines
2.3 KiB
Go
81 lines
2.3 KiB
Go
// Copyright 2026 Arctel.net
|
|
// SPDX-License-Identifier: Apache-2.0
|
|
|
|
// Package relay implements the relay node daemon runtime loop.
|
|
package relay
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"log/slog"
|
|
|
|
edgerunner "github.com/Rain-kl/Wavelet/internal/apps/edge/runner"
|
|
"github.com/Rain-kl/Wavelet/internal/apps/relay/config"
|
|
"github.com/Rain-kl/Wavelet/internal/apps/relay/frps"
|
|
"github.com/Rain-kl/Wavelet/internal/apps/relay/heartbeat"
|
|
"github.com/Rain-kl/Wavelet/internal/apps/relay/httpclient"
|
|
"github.com/Rain-kl/Wavelet/internal/apps/relay/state"
|
|
"github.com/Rain-kl/Wavelet/internal/apps/relay/wsclient"
|
|
service "github.com/Rain-kl/Wavelet/pkg/protocol"
|
|
)
|
|
|
|
// Runner manages the relay process.
|
|
type Runner struct {
|
|
Config *config.Config
|
|
StateStore *state.Store
|
|
HeartbeatService *heartbeat.Service
|
|
FrpsManager *frps.Manager
|
|
WebSocketService *wsclient.Client
|
|
HTTPClient *httpclient.Client
|
|
}
|
|
|
|
// Run starts the relay process by initiating the heartbeat and WS reconnection loop.
|
|
func (r *Runner) Run(ctx context.Context) error {
|
|
go r.HeartbeatService.Run(ctx)
|
|
|
|
return edgerunner.RunWSReconnectLoop(ctx, edgerunner.WSReconnectConfig{
|
|
ComponentName: "relay",
|
|
OnShutdown: r.FrpsManager.Stop,
|
|
}, func(ctx context.Context) (edgerunner.WSConnection, error) {
|
|
return r.WebSocketService.Connect(ctx)
|
|
}, func(ctx context.Context, conn edgerunner.WSConnection) {
|
|
r.handleConnection(ctx, conn)
|
|
})
|
|
}
|
|
|
|
type relayWSHandler struct {
|
|
runner *Runner
|
|
}
|
|
|
|
func (h *relayWSHandler) OnConnect(_ context.Context) error {
|
|
return nil
|
|
}
|
|
|
|
func (h *relayWSHandler) HandleMessage(ctx context.Context, msg wsclient.WSMessage) error {
|
|
switch msg.Type {
|
|
case "relay_config":
|
|
var cfg service.RelayConfig
|
|
if err := json.Unmarshal(msg.Payload, &cfg); err != nil {
|
|
slog.Error("failed to unmarshal relay_config", "error", err)
|
|
return nil
|
|
}
|
|
h.runner.FrpsManager.UpdateConfig(ctx, &cfg)
|
|
default:
|
|
slog.Debug("ignored unknown ws message type", "type", msg.Type)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (h *relayWSHandler) OnClose(err error) {
|
|
slog.Error("relay ws receive failed", "error", err)
|
|
}
|
|
|
|
func (r *Runner) handleConnection(ctx context.Context, conn edgerunner.WSConnection) {
|
|
wsConn, ok := conn.(*wsclient.Connection)
|
|
if !ok {
|
|
slog.Error("relay ws connection has unexpected type")
|
|
return
|
|
}
|
|
_ = wsConn.RunReceiveLoop(ctx, &relayWSHandler{runner: r})
|
|
}
|