mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-10-11 03:36:37 +08:00
fix(backend): resolve federation forward traffic stats and listener disappearance (#208)
* fix(backend): add repository methods for federation forward runtime management - GetActiveForwardPeerShareRuntimeByServiceName: lookup runtime by share_id and service_name - MarkForwardPeerShareRuntimeReleasedByServiceName: release runtime by service_name - ListActiveForwardPeerShareRuntimesByNodeAndServiceName: node-scoped query for flow processing Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai> * fix(backend): bind and release federation forward runtimes on service commands - bindPeerShareForwardRuntimeServices: create runtime if missing, update ServiceName/Port/Applied/Status - releasePeerShareForwardRuntimeServices: handle deleteservice command to mark runtime released - parseFederationForwardServiceNamesForRelease: extract service names from delete payload - Tests: bind creates runtime when missing, release marks runtime as released Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai> * fix(backend): scope federation flow lookup by node to avoid cross-share collisions - flowUpload: use GetNodeBySecret to extract nodeID for flow processing - processFlowItem: accept nodeID parameter and pass to flow handlers - processPeerShareFlowByServiceName: try node-scoped query first, fallback to global - Add warning log when multiple runtimes match (ambiguous) - Tests: update all processFlowItem calls with nodeID parameter Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai> * test(contract): adjust federation dual panel contract expectations Update assertion for entry share runtime binding behavior after fix Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai> --------- Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
@@ -1287,6 +1287,28 @@ func (r *Repository) ListActiveForwardPeerShareRuntimesByServiceName(serviceName
|
||||
return items, nil
|
||||
}
|
||||
|
||||
func (r *Repository) ListActiveForwardPeerShareRuntimesByNodeAndServiceName(nodeID int64, serviceName string) ([]model.PeerShareRuntime, error) {
|
||||
if r == nil || r.db == nil {
|
||||
return nil, errors.New("repository not initialized")
|
||||
}
|
||||
serviceName = strings.TrimSpace(serviceName)
|
||||
if serviceName == "" {
|
||||
return []model.PeerShareRuntime{}, nil
|
||||
}
|
||||
var items []model.PeerShareRuntime
|
||||
err := r.db.Where("node_id = ? AND service_name = ? AND status = 1 AND role = ?", nodeID, serviceName, "forward").
|
||||
Order("id ASC").
|
||||
Find(&items).Error
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if items == nil {
|
||||
items = make([]model.PeerShareRuntime, 0)
|
||||
}
|
||||
return items, nil
|
||||
}
|
||||
|
||||
|
||||
func (r *Repository) ListActiveForwardPeerShareRuntimeServiceNamesByNode(nodeID int64) ([]string, error) {
|
||||
if r == nil || r.db == nil {
|
||||
return nil, errors.New("repository not initialized")
|
||||
@@ -1333,6 +1355,27 @@ func (r *Repository) GetActiveForwardPeerShareRuntimeByPort(shareID int64, port
|
||||
return &item, nil
|
||||
}
|
||||
|
||||
func (r *Repository) GetActiveForwardPeerShareRuntimeByServiceName(shareID int64, serviceName string) (*model.PeerShareRuntime, error) {
|
||||
if r == nil || r.db == nil {
|
||||
return nil, errors.New("repository not initialized")
|
||||
}
|
||||
serviceName = strings.TrimSpace(serviceName)
|
||||
if shareID <= 0 || serviceName == "" {
|
||||
return nil, nil
|
||||
}
|
||||
var item model.PeerShareRuntime
|
||||
err := r.db.Where("share_id = ? AND service_name = ? AND status = 1 AND role = ?", shareID, serviceName, "forward").
|
||||
Order("id ASC").
|
||||
First(&item).Error
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil, nil
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &item, nil
|
||||
}
|
||||
|
||||
func (r *Repository) ExistsActivePeerShareRuntimeOnNodePort(nodeID int64, port int) (bool, error) {
|
||||
if r == nil || r.db == nil {
|
||||
return false, errors.New("repository not initialized")
|
||||
@@ -1376,6 +1419,27 @@ func (r *Repository) MarkPeerShareRuntimeReleasedByPort(shareID int64, port int,
|
||||
}).Error
|
||||
}
|
||||
|
||||
func (r *Repository) MarkForwardPeerShareRuntimeReleasedByServiceName(shareID int64, serviceName string, updatedTime int64) error {
|
||||
if r == nil || r.db == nil {
|
||||
return errors.New("repository not initialized")
|
||||
}
|
||||
serviceName = strings.TrimSpace(serviceName)
|
||||
if shareID <= 0 || serviceName == "" {
|
||||
return nil
|
||||
}
|
||||
if updatedTime <= 0 {
|
||||
updatedTime = unixMilliNow()
|
||||
}
|
||||
return r.db.Model(&model.PeerShareRuntime{}).
|
||||
Where("share_id = ? AND status = 1 AND role = ? AND service_name = ?", shareID, "forward", serviceName).
|
||||
Updates(map[string]interface{}{
|
||||
"status": 0,
|
||||
"applied": 0,
|
||||
"service_name": "",
|
||||
"updated_time": updatedTime,
|
||||
}).Error
|
||||
}
|
||||
|
||||
// ─── FederationTunnelBinding ─────────────────────────────────────────
|
||||
|
||||
func (r *Repository) UpsertFederationTunnelBinding(item *model.FederationTunnelBinding) error {
|
||||
|
||||
Reference in New Issue
Block a user