diff --git a/docker-compose-v4.yml b/docker-compose-v4.yml index ea85469..3e4bbe0 100644 --- a/docker-compose-v4.yml +++ b/docker-compose-v4.yml @@ -87,4 +87,4 @@ networks: driver: bridge ipam: config: - - subnet: 172.20.0.0/16 + - subnet: 172.80.0.0/16 diff --git a/docker-compose-v6.yml b/docker-compose-v6.yml index cb7ae47..b77fc2f 100644 --- a/docker-compose-v6.yml +++ b/docker-compose-v6.yml @@ -88,5 +88,5 @@ networks: enable_ipv6: true ipam: config: - - subnet: 172.20.0.0/16 + - subnet: 172.80.0.0/16 - subnet: fd00:dead:beef::/48 diff --git a/go-backend/internal/http/client/federation.go b/go-backend/internal/http/client/federation.go index 128055a..09d6a64 100644 --- a/go-backend/internal/http/client/federation.go +++ b/go-backend/internal/http/client/federation.go @@ -77,6 +77,18 @@ type RuntimeDiagnoseRequest struct { 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 { return &FederationClient{ client: &http.Client{ @@ -333,3 +345,42 @@ func (c *FederationClient) Diagnose(url, token, localDomain string, reqData Runt 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 +} diff --git a/go-backend/internal/http/handler/control_plane.go b/go-backend/internal/http/handler/control_plane.go index 0e3d084..d61a86a 100644 --- a/go-backend/internal/http/handler/control_plane.go +++ b/go-backend/internal/http/handler/control_plane.go @@ -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) { - 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 { return result, nil } @@ -475,6 +485,44 @@ func (h *Handler) sendNodeCommand(nodeID int64, commandType string, data interfa 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) { if forward == nil { return nil, errForwardNotFound diff --git a/go-backend/internal/http/handler/federation.go b/go-backend/internal/http/handler/federation.go index 09abef2..831f64f 100644 --- a/go-backend/internal/http/handler/federation.go +++ b/go-backend/internal/http/handler/federation.go @@ -91,6 +91,11 @@ type federationRuntimeDiagnoseRequest struct { Timeout int `json:"timeout"` } +type federationRuntimeCommandRequest struct { + CommandType string `json:"commandType"` + Data interface{} `json:"data"` +} + type peerShareUsedPort struct { RuntimeID int64 `json:"runtimeId"` 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)) } +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) { if share == nil { return 0, fmt.Errorf("share not found") diff --git a/go-backend/internal/http/handler/federation_runtime_test.go b/go-backend/internal/http/handler/federation_runtime_test.go index b29f758..17c1520 100644 --- a/go-backend/internal/http/handler/federation_runtime_test.go +++ b/go-backend/internal/http/handler/federation_runtime_test.go @@ -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) { repo, err := sqlite.Open(filepath.Join(t.TempDir(), "panel.db")) if err != nil { diff --git a/go-backend/internal/http/handler/handler.go b/go-backend/internal/http/handler/handler.go index 8a26877..3bd4170 100644 --- a/go-backend/internal/http/handler/handler.go +++ b/go-backend/internal/http/handler/handler.go @@ -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/release-role", h.authPeer(h.federationRuntimeReleaseRole)) 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/backup/export", h.backupExport) diff --git a/go-backend/internal/http/handler/mutations.go b/go-backend/internal/http/handler/mutations.go index f05f66e..c2cc211 100644 --- a/go-backend/internal/http/handler/mutations.go +++ b/go-backend/internal/http/handler/mutations.go @@ -2233,7 +2233,7 @@ func (h *Handler) prepareTunnelCreateState(tx *store.Tx, req map[string]interfac } return nil, err } - if node.Status != 1 { + if node.IsRemote != 1 && node.Status != 1 { return nil, errors.New("部分节点不在线") } state.Nodes[nodeID] = node diff --git a/go-backend/internal/http/middleware/auth.go b/go-backend/internal/http/middleware/auth.go index bf3aebf..7e60afc 100644 --- a/go-backend/internal/http/middleware/auth.go +++ b/go-backend/internal/http/middleware/auth.go @@ -93,6 +93,8 @@ func shouldSkip(path string) bool { return true case path == "/api/v1/federation/runtime/diagnose": return true + case path == "/api/v1/federation/runtime/command": + return true default: return false } diff --git a/go-backend/tests/contract/federation_dual_panel_contract_test.go b/go-backend/tests/contract/federation_dual_panel_contract_test.go index b3bece9..0959f65 100644 --- a/go-backend/tests/contract/federation_dual_panel_contract_test.go +++ b/go-backend/tests/contract/federation_dual_panel_contract_test.go @@ -78,6 +78,8 @@ func TestFederationDualPanelMiddleExitAutoPortContract(t *testing.T) { middleRemoteNodeID := queryRemoteNodeIDByToken(t, consumerRepo, "share-middle-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") defer stopMiddle() 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, 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`, exitShareID, 1) assertCount(t, providerRepo, `SELECT COUNT(1) FROM peer_share_runtime WHERE share_id = ?`, entryShareID, 0)