mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-10-03 17:06:36 +08:00
feat(federation): add remote node command support
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
@@ -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")
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user