mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-10-06 10:06:36 +08:00
Merge branch 'main' into opencode/quick-comet
This commit is contained in:
@@ -87,4 +87,4 @@ networks:
|
|||||||
driver: bridge
|
driver: bridge
|
||||||
ipam:
|
ipam:
|
||||||
config:
|
config:
|
||||||
- subnet: 172.20.0.0/16
|
- subnet: 172.80.0.0/16
|
||||||
|
|||||||
@@ -88,5 +88,5 @@ networks:
|
|||||||
enable_ipv6: true
|
enable_ipv6: true
|
||||||
ipam:
|
ipam:
|
||||||
config:
|
config:
|
||||||
- subnet: 172.20.0.0/16
|
- subnet: 172.80.0.0/16
|
||||||
- subnet: fd00:dead:beef::/48
|
- subnet: fd00:dead:beef::/48
|
||||||
|
|||||||
@@ -77,6 +77,18 @@ type RuntimeDiagnoseRequest struct {
|
|||||||
Timeout int `json:"timeout"`
|
Timeout int `json:"timeout"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type RuntimeNodeCommandRequest struct {
|
||||||
|
CommandType string `json:"commandType"`
|
||||||
|
Data interface{} `json:"data"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type RuntimeNodeCommandResponse struct {
|
||||||
|
Type string `json:"type"`
|
||||||
|
Success bool `json:"success"`
|
||||||
|
Message string `json:"message"`
|
||||||
|
Data map[string]interface{} `json:"data,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
func NewFederationClient() *FederationClient {
|
func NewFederationClient() *FederationClient {
|
||||||
return &FederationClient{
|
return &FederationClient{
|
||||||
client: &http.Client{
|
client: &http.Client{
|
||||||
@@ -333,3 +345,42 @@ func (c *FederationClient) Diagnose(url, token, localDomain string, reqData Runt
|
|||||||
|
|
||||||
return res.Data, nil
|
return res.Data, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (c *FederationClient) Command(url, token, localDomain string, reqData RuntimeNodeCommandRequest) (*RuntimeNodeCommandResponse, error) {
|
||||||
|
url = strings.TrimSuffix(url, "/")
|
||||||
|
bodyBytes, _ := json.Marshal(reqData)
|
||||||
|
req, err := http.NewRequest("POST", url+"/api/v1/federation/runtime/command", strings.NewReader(string(bodyBytes)))
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
req.Header.Set("Authorization", "Bearer "+token)
|
||||||
|
if localDomain != "" {
|
||||||
|
req.Header.Set("X-Panel-Domain", localDomain)
|
||||||
|
}
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
|
||||||
|
resp, err := c.client.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
defer resp.Body.Close()
|
||||||
|
|
||||||
|
if resp.StatusCode != 200 {
|
||||||
|
body, _ := io.ReadAll(resp.Body)
|
||||||
|
return nil, fmt.Errorf("remote error %d: %s", resp.StatusCode, string(body))
|
||||||
|
}
|
||||||
|
|
||||||
|
var res struct {
|
||||||
|
Code int `json:"code"`
|
||||||
|
Msg string `json:"msg"`
|
||||||
|
Data RuntimeNodeCommandResponse `json:"data"`
|
||||||
|
}
|
||||||
|
if err := json.NewDecoder(resp.Body).Decode(&res); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if res.Code != 0 {
|
||||||
|
return nil, fmt.Errorf("remote api error: %s", res.Msg)
|
||||||
|
}
|
||||||
|
|
||||||
|
return &res.Data, nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -457,7 +457,17 @@ func (h *Handler) applyNodeProtocolChange(nodeID int64, httpVal, tlsVal, socksVa
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (h *Handler) sendNodeCommand(nodeID int64, commandType string, data interface{}, tolerateExists bool, tolerateNotFound bool) (ws.CommandResult, error) {
|
func (h *Handler) sendNodeCommand(nodeID int64, commandType string, data interface{}, tolerateExists bool, tolerateNotFound bool) (ws.CommandResult, error) {
|
||||||
result, err := h.wsServer.SendCommand(nodeID, commandType, data, 12*time.Second)
|
var (
|
||||||
|
result ws.CommandResult
|
||||||
|
err error
|
||||||
|
)
|
||||||
|
|
||||||
|
node, nodeErr := h.getNodeRecord(nodeID)
|
||||||
|
if nodeErr == nil && node != nil && node.IsRemote == 1 {
|
||||||
|
result, err = h.sendRemoteNodeCommand(node, commandType, data)
|
||||||
|
} else {
|
||||||
|
result, err = h.wsServer.SendCommand(nodeID, commandType, data, 12*time.Second)
|
||||||
|
}
|
||||||
if err == nil {
|
if err == nil {
|
||||||
return result, nil
|
return result, nil
|
||||||
}
|
}
|
||||||
@@ -475,6 +485,44 @@ func (h *Handler) sendNodeCommand(nodeID int64, commandType string, data interfa
|
|||||||
return result, err
|
return result, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (h *Handler) sendRemoteNodeCommand(node *nodeRecord, commandType string, data interface{}) (ws.CommandResult, error) {
|
||||||
|
if node == nil {
|
||||||
|
return ws.CommandResult{}, errors.New("节点不存在")
|
||||||
|
}
|
||||||
|
remoteURL := strings.TrimSpace(node.RemoteURL)
|
||||||
|
remoteToken := strings.TrimSpace(node.RemoteToken)
|
||||||
|
if remoteURL == "" || remoteToken == "" {
|
||||||
|
return ws.CommandResult{}, errors.New("远程节点缺少共享配置")
|
||||||
|
}
|
||||||
|
|
||||||
|
fc := client.NewFederationClient()
|
||||||
|
res, err := fc.Command(remoteURL, remoteToken, h.federationLocalDomain(), client.RuntimeNodeCommandRequest{
|
||||||
|
CommandType: commandType,
|
||||||
|
Data: data,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
return ws.CommandResult{}, err
|
||||||
|
}
|
||||||
|
if res == nil {
|
||||||
|
return ws.CommandResult{}, errors.New("远程节点未返回命令结果")
|
||||||
|
}
|
||||||
|
|
||||||
|
result := ws.CommandResult{
|
||||||
|
Type: res.Type,
|
||||||
|
Success: res.Success,
|
||||||
|
Message: res.Message,
|
||||||
|
Data: res.Data,
|
||||||
|
}
|
||||||
|
if !result.Success {
|
||||||
|
msg := strings.TrimSpace(result.Message)
|
||||||
|
if msg == "" {
|
||||||
|
msg = "命令执行失败"
|
||||||
|
}
|
||||||
|
return result, errors.New(msg)
|
||||||
|
}
|
||||||
|
return result, nil
|
||||||
|
}
|
||||||
|
|
||||||
func (h *Handler) diagnoseForwardRuntime(forward *forwardRecord) (map[string]interface{}, error) {
|
func (h *Handler) diagnoseForwardRuntime(forward *forwardRecord) (map[string]interface{}, error) {
|
||||||
if forward == nil {
|
if forward == nil {
|
||||||
return nil, errForwardNotFound
|
return nil, errForwardNotFound
|
||||||
|
|||||||
@@ -91,6 +91,11 @@ type federationRuntimeDiagnoseRequest struct {
|
|||||||
Timeout int `json:"timeout"`
|
Timeout int `json:"timeout"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type federationRuntimeCommandRequest struct {
|
||||||
|
CommandType string `json:"commandType"`
|
||||||
|
Data interface{} `json:"data"`
|
||||||
|
}
|
||||||
|
|
||||||
type peerShareUsedPort struct {
|
type peerShareUsedPort struct {
|
||||||
RuntimeID int64 `json:"runtimeId"`
|
RuntimeID int64 `json:"runtimeId"`
|
||||||
Port int `json:"port"`
|
Port int `json:"port"`
|
||||||
@@ -1199,6 +1204,51 @@ func (h *Handler) federationRuntimeDiagnose(w http.ResponseWriter, r *http.Reque
|
|||||||
response.WriteJSON(w, response.OK(res.Data))
|
response.WriteJSON(w, response.OK(res.Data))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (h *Handler) federationRuntimeCommand(w http.ResponseWriter, r *http.Request) {
|
||||||
|
if r.Method != http.MethodPost {
|
||||||
|
response.WriteJSON(w, response.ErrDefault("Invalid method"))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
token := extractBearerToken(r)
|
||||||
|
share, err := h.repo.GetPeerShareByToken(token)
|
||||||
|
if err != nil || share == nil {
|
||||||
|
response.WriteJSON(w, response.Err(401, "Unauthorized"))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
var req federationRuntimeCommandRequest
|
||||||
|
if err := decodeJSON(r.Body, &req); err != nil {
|
||||||
|
response.WriteJSON(w, response.ErrDefault("Invalid JSON"))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
cmd := strings.TrimSpace(req.CommandType)
|
||||||
|
if cmd == "" {
|
||||||
|
response.WriteJSON(w, response.ErrDefault("commandType is required"))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if !isFederationRuntimeCommandAllowed(cmd) {
|
||||||
|
response.WriteJSON(w, response.ErrDefault("command not allowed"))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
res, err := h.sendNodeCommand(share.NodeID, cmd, req.Data, false, false)
|
||||||
|
if err != nil {
|
||||||
|
response.WriteJSON(w, response.ErrDefault(err.Error()))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
response.WriteJSON(w, response.OK(res))
|
||||||
|
}
|
||||||
|
|
||||||
|
func isFederationRuntimeCommandAllowed(commandType string) bool {
|
||||||
|
switch strings.ToLower(strings.TrimSpace(commandType)) {
|
||||||
|
case "addservice", "updateservice", "deleteservice", "pauseservice", "resumeservice", "addchains", "deletechains", "addlimiters", "deletelimiters", "tcpping", "reload":
|
||||||
|
return true
|
||||||
|
default:
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (h *Handler) pickPeerSharePort(share *sqlite.PeerShare, requestedPort int) (int, error) {
|
func (h *Handler) pickPeerSharePort(share *sqlite.PeerShare, requestedPort int) (int, error) {
|
||||||
if share == nil {
|
if share == nil {
|
||||||
return 0, fmt.Errorf("share not found")
|
return 0, fmt.Errorf("share not found")
|
||||||
|
|||||||
@@ -152,6 +152,71 @@ func TestPrepareTunnelCreateStateRemoteAutoPortDefersToFederation(t *testing.T)
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestPrepareTunnelCreateStateAllowsOfflineRemoteMiddleNode(t *testing.T) {
|
||||||
|
repo, err := sqlite.Open(filepath.Join(t.TempDir(), "panel.db"))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("open repo: %v", err)
|
||||||
|
}
|
||||||
|
defer repo.Close()
|
||||||
|
|
||||||
|
h := &Handler{repo: repo}
|
||||||
|
now := time.Now().UnixMilli()
|
||||||
|
|
||||||
|
insertNode := func(name string, status int, portRange string, isRemote int) int64 {
|
||||||
|
res, execErr := repo.DB().Exec(`
|
||||||
|
INSERT INTO node(name, secret, server_ip, server_ip_v4, server_ip_v6, port, interface_name, version, http, tls, socks, created_time, updated_time, status, tcp_listen_addr, udp_listen_addr, inx, is_remote, remote_url, remote_token, remote_config)
|
||||||
|
VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||||
|
`, name, name+"-secret", "10.0.0.1", "10.0.0.1", "", portRange, "", "v1", 1, 1, 1, now, now, status, "[::]", "[::]", 0, isRemote, "http://peer", "peer-token", `{"shareId":2}`)
|
||||||
|
if execErr != nil {
|
||||||
|
t.Fatalf("insert node %s: %v", name, execErr)
|
||||||
|
}
|
||||||
|
id, idErr := res.LastInsertId()
|
||||||
|
if idErr != nil {
|
||||||
|
t.Fatalf("node id %s: %v", name, idErr)
|
||||||
|
}
|
||||||
|
return id
|
||||||
|
}
|
||||||
|
|
||||||
|
entryID := insertNode("entry-local", 1, "32000-32010", 0)
|
||||||
|
remoteMiddleID := insertNode("middle-remote", 0, "33000-33010", 1)
|
||||||
|
outID := insertNode("out-local", 1, "34000-34010", 0)
|
||||||
|
|
||||||
|
tx, err := repo.DB().Begin()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("begin tx: %v", err)
|
||||||
|
}
|
||||||
|
defer tx.Rollback()
|
||||||
|
|
||||||
|
req := map[string]interface{}{
|
||||||
|
"name": "remote-middle-offline-status",
|
||||||
|
"inNodeId": []interface{}{
|
||||||
|
map[string]interface{}{"nodeId": float64(entryID), "protocol": "tls", "strategy": "round"},
|
||||||
|
},
|
||||||
|
"chainNodes": []interface{}{
|
||||||
|
[]interface{}{
|
||||||
|
map[string]interface{}{"nodeId": float64(remoteMiddleID), "protocol": "tls", "strategy": "round", "port": float64(0)},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
"outNodeId": []interface{}{
|
||||||
|
map[string]interface{}{"nodeId": float64(outID), "protocol": "tls", "strategy": "round", "port": float64(0)},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
state, err := h.prepareTunnelCreateState(tx, req, 2, 0)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("prepare state should allow offline remote middle node: %v", err)
|
||||||
|
}
|
||||||
|
if len(state.ChainHops) != 1 || len(state.ChainHops[0]) != 1 {
|
||||||
|
t.Fatalf("expected one middle hop node, got %+v", state.ChainHops)
|
||||||
|
}
|
||||||
|
if state.ChainHops[0][0].NodeID != remoteMiddleID {
|
||||||
|
t.Fatalf("expected remote middle node id %d, got %d", remoteMiddleID, state.ChainHops[0][0].NodeID)
|
||||||
|
}
|
||||||
|
if state.Nodes[remoteMiddleID] == nil || state.Nodes[remoteMiddleID].IsRemote != 1 {
|
||||||
|
t.Fatalf("expected remote middle node metadata in state")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestFederationRuntimeReservePortRejectsWhenShareFlowExceeded(t *testing.T) {
|
func TestFederationRuntimeReservePortRejectsWhenShareFlowExceeded(t *testing.T) {
|
||||||
repo, err := sqlite.Open(filepath.Join(t.TempDir(), "panel.db"))
|
repo, err := sqlite.Open(filepath.Join(t.TempDir(), "panel.db"))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -169,6 +169,7 @@ func (h *Handler) Register(mux *http.ServeMux) {
|
|||||||
mux.HandleFunc("/api/v1/federation/runtime/apply-role", h.authPeer(h.federationRuntimeApplyRole))
|
mux.HandleFunc("/api/v1/federation/runtime/apply-role", h.authPeer(h.federationRuntimeApplyRole))
|
||||||
mux.HandleFunc("/api/v1/federation/runtime/release-role", h.authPeer(h.federationRuntimeReleaseRole))
|
mux.HandleFunc("/api/v1/federation/runtime/release-role", h.authPeer(h.federationRuntimeReleaseRole))
|
||||||
mux.HandleFunc("/api/v1/federation/runtime/diagnose", h.authPeer(h.federationRuntimeDiagnose))
|
mux.HandleFunc("/api/v1/federation/runtime/diagnose", h.authPeer(h.federationRuntimeDiagnose))
|
||||||
|
mux.HandleFunc("/api/v1/federation/runtime/command", h.authPeer(h.federationRuntimeCommand))
|
||||||
mux.HandleFunc("/api/v1/federation/node/import", h.nodeImport)
|
mux.HandleFunc("/api/v1/federation/node/import", h.nodeImport)
|
||||||
|
|
||||||
mux.HandleFunc("/api/v1/backup/export", h.backupExport)
|
mux.HandleFunc("/api/v1/backup/export", h.backupExport)
|
||||||
|
|||||||
@@ -2233,7 +2233,7 @@ func (h *Handler) prepareTunnelCreateState(tx *store.Tx, req map[string]interfac
|
|||||||
}
|
}
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
if node.Status != 1 {
|
if node.IsRemote != 1 && node.Status != 1 {
|
||||||
return nil, errors.New("部分节点不在线")
|
return nil, errors.New("部分节点不在线")
|
||||||
}
|
}
|
||||||
state.Nodes[nodeID] = node
|
state.Nodes[nodeID] = node
|
||||||
|
|||||||
@@ -93,6 +93,8 @@ func shouldSkip(path string) bool {
|
|||||||
return true
|
return true
|
||||||
case path == "/api/v1/federation/runtime/diagnose":
|
case path == "/api/v1/federation/runtime/diagnose":
|
||||||
return true
|
return true
|
||||||
|
case path == "/api/v1/federation/runtime/command":
|
||||||
|
return true
|
||||||
default:
|
default:
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -78,6 +78,8 @@ func TestFederationDualPanelMiddleExitAutoPortContract(t *testing.T) {
|
|||||||
middleRemoteNodeID := queryRemoteNodeIDByToken(t, consumerRepo, "share-middle-token")
|
middleRemoteNodeID := queryRemoteNodeIDByToken(t, consumerRepo, "share-middle-token")
|
||||||
exitRemoteNodeID := queryRemoteNodeIDByToken(t, consumerRepo, "share-exit-token")
|
exitRemoteNodeID := queryRemoteNodeIDByToken(t, consumerRepo, "share-exit-token")
|
||||||
|
|
||||||
|
stopEntry := startMockNodeSession(t, providerServer.URL, "provider-entry-secret")
|
||||||
|
defer stopEntry()
|
||||||
stopMiddle := startMockNodeSession(t, providerServer.URL, "provider-middle-secret")
|
stopMiddle := startMockNodeSession(t, providerServer.URL, "provider-middle-secret")
|
||||||
defer stopMiddle()
|
defer stopMiddle()
|
||||||
stopExit := startMockNodeSession(t, providerServer.URL, "provider-exit-secret")
|
stopExit := startMockNodeSession(t, providerServer.URL, "provider-exit-secret")
|
||||||
@@ -148,6 +150,23 @@ func TestFederationDualPanelMiddleExitAutoPortContract(t *testing.T) {
|
|||||||
assertTunnelPortInRange(t, consumerRepo, secondTunnelID, 2, middleRemoteNodeID, 44000, 44010)
|
assertTunnelPortInRange(t, consumerRepo, secondTunnelID, 2, middleRemoteNodeID, 44000, 44010)
|
||||||
assertTunnelPortInRange(t, consumerRepo, secondTunnelID, 3, exitRemoteNodeID, 45000, 45010)
|
assertTunnelPortInRange(t, consumerRepo, secondTunnelID, 3, exitRemoteNodeID, 45000, 45010)
|
||||||
|
|
||||||
|
forwardPayload := map[string]interface{}{
|
||||||
|
"name": "dual-panel-remote-entry-forward",
|
||||||
|
"tunnelId": secondTunnelID,
|
||||||
|
"remoteAddr": "1.1.1.1:443",
|
||||||
|
"strategy": "fifo",
|
||||||
|
}
|
||||||
|
forwardBody, err := json.Marshal(forwardPayload)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("marshal forward payload: %v", err)
|
||||||
|
}
|
||||||
|
forwardReq := httptest.NewRequest(http.MethodPost, "/api/v1/forward/create", bytes.NewReader(forwardBody))
|
||||||
|
forwardReq.Header.Set("Authorization", consumerAdminToken)
|
||||||
|
forwardReq.Header.Set("Content-Type", "application/json")
|
||||||
|
forwardRes := httptest.NewRecorder()
|
||||||
|
consumerRouter.ServeHTTP(forwardRes, forwardReq)
|
||||||
|
assertCode(t, forwardRes, 0)
|
||||||
|
|
||||||
assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ? AND status = 1 AND applied = 1`, middleShareID, 1)
|
assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ? AND status = 1 AND applied = 1`, middleShareID, 1)
|
||||||
assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ? AND status = 1 AND applied = 1`, exitShareID, 1)
|
assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ? AND status = 1 AND applied = 1`, exitShareID, 1)
|
||||||
assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ?`, entryShareID, 0)
|
assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ?`, entryShareID, 0)
|
||||||
|
|||||||
Reference in New Issue
Block a user