Merge pull request #104 from Sagit-chu/opencode/tidy-cactus

fix: 修复共享节点作为入口的问题和调整docker 网络
This commit is contained in:
sagit
2026-02-13 14:16:56 +08:00
committed by GitHub
10 changed files with 240 additions and 4 deletions
+1 -1
View File
@@ -87,4 +87,4 @@ networks:
driver: bridge
ipam:
config:
- subnet: 172.20.0.0/16
- subnet: 172.80.0.0/16
+1 -1
View File
@@ -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
@@ -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
}
@@ -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
@@ -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")
@@ -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 {
@@ -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("/flow/test", h.flowTest)
@@ -2061,7 +2061,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
@@ -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
}
@@ -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)