mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-10-06 18:06:36 +08:00
fix(federation): sync remote node status on list and detect deleted provider shares
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
@@ -85,6 +85,14 @@ func NewFederationClient() *FederationClient {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func NewFederationClientWithTimeout(timeout time.Duration) *FederationClient {
|
||||||
|
return &FederationClient{
|
||||||
|
client: &http.Client{
|
||||||
|
Timeout: timeout,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (c *FederationClient) Connect(url, token, localDomain string) (*RemoteNodeInfo, error) {
|
func (c *FederationClient) Connect(url, token, localDomain string) (*RemoteNodeInfo, error) {
|
||||||
url = strings.TrimSuffix(url, "/")
|
url = strings.TrimSuffix(url, "/")
|
||||||
req, err := http.NewRequest("POST", url+"/api/v1/federation/connect", nil)
|
req, err := http.NewRequest("POST", url+"/api/v1/federation/connect", nil)
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"sort"
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"go-backend/internal/http/client"
|
"go-backend/internal/http/client"
|
||||||
@@ -1402,6 +1403,74 @@ func isPeerIPAllowed(clientIP net.IP, whitelist string) bool {
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (h *Handler) syncRemoteNodeStatuses(items []map[string]interface{}) {
|
||||||
|
type remoteEntry struct {
|
||||||
|
index int
|
||||||
|
remoteURL string
|
||||||
|
remoteToken string
|
||||||
|
}
|
||||||
|
|
||||||
|
var remotes []remoteEntry
|
||||||
|
for i, item := range items {
|
||||||
|
isRemote, _ := item["isRemote"].(int)
|
||||||
|
if isRemote != 1 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
url, _ := item["remoteUrl"].(string)
|
||||||
|
token, _ := item["remoteToken"].(string)
|
||||||
|
url = strings.TrimSpace(url)
|
||||||
|
token = strings.TrimSpace(token)
|
||||||
|
if url == "" || token == "" {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
remotes = append(remotes, remoteEntry{index: i, remoteURL: url, remoteToken: token})
|
||||||
|
}
|
||||||
|
if len(remotes) == 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
localDomain := h.federationLocalDomain()
|
||||||
|
fc := client.NewFederationClientWithTimeout(5 * time.Second)
|
||||||
|
|
||||||
|
type syncResult struct {
|
||||||
|
index int
|
||||||
|
status int
|
||||||
|
syncError string
|
||||||
|
}
|
||||||
|
|
||||||
|
results := make([]syncResult, len(remotes))
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
for i, entry := range remotes {
|
||||||
|
wg.Add(1)
|
||||||
|
go func(idx int, e remoteEntry) {
|
||||||
|
defer wg.Done()
|
||||||
|
info, err := fc.Connect(e.remoteURL, e.remoteToken, localDomain)
|
||||||
|
if err != nil {
|
||||||
|
errMsg := err.Error()
|
||||||
|
if strings.Contains(errMsg, "401") || strings.Contains(errMsg, "Invalid token") || strings.Contains(errMsg, "Unauthorized") {
|
||||||
|
results[idx] = syncResult{index: e.index, status: 0, syncError: "provider_share_deleted"}
|
||||||
|
} else if strings.Contains(errMsg, "403") || strings.Contains(errMsg, "Share is disabled") {
|
||||||
|
results[idx] = syncResult{index: e.index, status: 0, syncError: "provider_share_disabled"}
|
||||||
|
} else if strings.Contains(errMsg, "Share expired") {
|
||||||
|
results[idx] = syncResult{index: e.index, status: 0, syncError: "provider_share_expired"}
|
||||||
|
} else {
|
||||||
|
results[idx] = syncResult{index: e.index, status: 0, syncError: errMsg}
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
results[idx] = syncResult{index: e.index, status: info.Status, syncError: ""}
|
||||||
|
}
|
||||||
|
}(i, entry)
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
|
|
||||||
|
for _, r := range results {
|
||||||
|
items[r.index]["status"] = r.status
|
||||||
|
if r.syncError != "" {
|
||||||
|
items[r.index]["syncError"] = r.syncError
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (h *Handler) cleanupPeerShareRuntimes(shareID int64) {
|
func (h *Handler) cleanupPeerShareRuntimes(shareID int64) {
|
||||||
if h == nil || h.repo == nil || shareID <= 0 {
|
if h == nil || h.repo == nil || shareID <= 0 {
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -321,6 +321,9 @@ func (h *Handler) nodeList(w http.ResponseWriter, r *http.Request) {
|
|||||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
h.syncRemoteNodeStatuses(items)
|
||||||
|
|
||||||
response.WriteJSON(w, response.OK(items))
|
response.WriteJSON(w, response.OK(items))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user