feat(api): enhance peer session management and documentation
CI / changes (push) Successful in 8s
CI / commitlint (push) Has been skipped
CI / openapi (push) Successful in 27s
CI / web (push) Failing after 35s
CI / go (push) Successful in 55s
CI / bird2 (push) Successful in 15s
CI / release (push) Has been skipped
CI / changes (push) Successful in 8s
CI / commitlint (push) Has been skipped
CI / openapi (push) Successful in 27s
CI / web (push) Failing after 35s
CI / go (push) Successful in 55s
CI / bird2 (push) Successful in 15s
CI / release (push) Has been skipped
- Added new fields to the API for tracking connected speakers and session states across multiple nodes, including `connected_speaker_id`, `connected_speaker_label`, `session_on_speakers`, `established_on_speakers`, and `session_mismatch`. - Implemented a new endpoint for retrieving bird protocol sessions, enhancing the agent server functionality. - Updated the OpenAPI documentation to reflect the new fields and query parameters, improving clarity for API consumers. - Modified the frontend to display connected speaker information and session states, providing better visibility into peer connections. - Updated deployment documentation to clarify the configuration requirements for enabling IP forwarding on VPS.
This commit is contained in:
@@ -0,0 +1,263 @@
|
||||
package httpapi
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/netip"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"evobgp/internal/birdfmt"
|
||||
"evobgp/internal/nodedispatch"
|
||||
"evobgp/internal/store"
|
||||
)
|
||||
|
||||
const peerLiveCacheTTL = 15 * time.Second
|
||||
|
||||
type speakerBGPLive struct {
|
||||
SpeakerID string
|
||||
Label string
|
||||
Sessions []birdfmt.BGPSession
|
||||
Error string
|
||||
}
|
||||
|
||||
type peerLiveCacheEntry struct {
|
||||
at time.Time
|
||||
views []speakerBGPLive
|
||||
}
|
||||
|
||||
var peerLiveCache sync.Map // tenantID -> peerLiveCacheEntry
|
||||
|
||||
type peerSessionOnSpeaker struct {
|
||||
SpeakerID string `json:"speaker_id"`
|
||||
Label string `json:"label"`
|
||||
State string `json:"state"`
|
||||
}
|
||||
|
||||
func speakerDisplayLabel(sp *store.Speaker) string {
|
||||
if sp == nil {
|
||||
return ""
|
||||
}
|
||||
meta := store.ParseSpeakerMeta(sp.MetaJSON)
|
||||
host := strings.TrimSpace(meta.AgentDomain)
|
||||
if host == "" {
|
||||
host = strings.TrimSpace(sp.Endpoint)
|
||||
}
|
||||
if strings.EqualFold(strings.TrimSpace(sp.Role), "master") {
|
||||
if host != "" {
|
||||
return "CP · " + host
|
||||
}
|
||||
return "CP (master)"
|
||||
}
|
||||
if host != "" {
|
||||
return host
|
||||
}
|
||||
return sp.ID
|
||||
}
|
||||
|
||||
func masterSpeakerID(speakers []*store.Speaker) string {
|
||||
for _, sp := range speakers {
|
||||
if sp != nil && strings.EqualFold(strings.TrimSpace(sp.Role), "master") {
|
||||
return sp.ID
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func (s *Server) collectSpeakerBGPLive(ctx context.Context, tenantID string, fresh bool) []speakerBGPLive {
|
||||
if !fresh {
|
||||
if v, ok := peerLiveCache.Load(tenantID); ok {
|
||||
ent := v.(peerLiveCacheEntry)
|
||||
if time.Since(ent.at) < peerLiveCacheTTL {
|
||||
return ent.views
|
||||
}
|
||||
}
|
||||
}
|
||||
speakers := s.store.ListSpeakersForTenant(tenantID)
|
||||
views := make([]speakerBGPLive, 0, len(speakers)+1)
|
||||
|
||||
if sock := strings.TrimSpace(os.Getenv("EVOBGP_BIRDC_SOCKET")); sock != "" {
|
||||
v := speakerBGPLive{Label: "CP (local BIRD)"}
|
||||
if mid := masterSpeakerID(speakers); mid != "" {
|
||||
v.SpeakerID = mid
|
||||
for _, sp := range speakers {
|
||||
if sp != nil && sp.ID == mid {
|
||||
v.Label = speakerDisplayLabel(sp)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
out, err := birdfmt.ShowProtocols(ctx, sock, strings.TrimSpace(os.Getenv("EVOBGP_BIRDC_BIN")))
|
||||
if err != nil {
|
||||
v.Error = err.Error()
|
||||
} else {
|
||||
v.Sessions = birdfmt.ParseBGPSessions(out)
|
||||
}
|
||||
views = append(views, v)
|
||||
}
|
||||
|
||||
opts := nodedispatch.Options{Timeout: 8 * time.Second}
|
||||
type resWrap struct {
|
||||
sp *store.Speaker
|
||||
res nodedispatch.BirdProtocolsResult
|
||||
}
|
||||
ch := make(chan resWrap, len(speakers))
|
||||
var wg sync.WaitGroup
|
||||
for _, sp := range speakers {
|
||||
if sp == nil {
|
||||
continue
|
||||
}
|
||||
meta := store.ParseSpeakerMeta(sp.MetaJSON)
|
||||
if !store.SpeakerNeedsRemoteDispatch(sp.Role, meta) {
|
||||
continue
|
||||
}
|
||||
wg.Add(1)
|
||||
go func(speaker *store.Speaker) {
|
||||
defer wg.Done()
|
||||
ch <- resWrap{
|
||||
sp: speaker,
|
||||
res: nodedispatch.FetchBirdProtocols(ctx, speaker, opts),
|
||||
}
|
||||
}(sp)
|
||||
}
|
||||
wg.Wait()
|
||||
close(ch)
|
||||
for rw := range ch {
|
||||
views = append(views, speakerBGPLive{
|
||||
SpeakerID: rw.sp.ID,
|
||||
Label: speakerDisplayLabel(rw.sp),
|
||||
Sessions: rw.res.Sessions,
|
||||
Error: rw.res.Error,
|
||||
})
|
||||
}
|
||||
|
||||
peerLiveCache.Store(tenantID, peerLiveCacheEntry{at: time.Now(), views: views})
|
||||
return views
|
||||
}
|
||||
|
||||
func matchPeerOnSpeakers(peer *store.BGPPeer, views []speakerBGPLive) (
|
||||
bestState string,
|
||||
connectedID string,
|
||||
connectedLabel string,
|
||||
establishedOn []peerSessionOnSpeaker,
|
||||
on []peerSessionOnSpeaker,
|
||||
mismatch bool,
|
||||
) {
|
||||
if peer == nil {
|
||||
return "", "", "", nil, nil, false
|
||||
}
|
||||
neighbor, hasNeighbor := store.ParsePeerNeighbor(peer.Neighbor)
|
||||
protoName := birdfmt.PeerProtocolName(peer.ID)
|
||||
|
||||
for _, v := range views {
|
||||
if v.Error != "" && len(v.Sessions) == 0 {
|
||||
continue
|
||||
}
|
||||
for _, sess := range v.Sessions {
|
||||
if !peerSessionMatches(sess, protoName, neighbor, hasNeighbor) {
|
||||
continue
|
||||
}
|
||||
hit := peerSessionOnSpeaker{
|
||||
SpeakerID: v.SpeakerID,
|
||||
Label: v.Label,
|
||||
State: sess.State,
|
||||
}
|
||||
on = append(on, hit)
|
||||
if strings.EqualFold(strings.TrimSpace(sess.State), "Established") {
|
||||
establishedOn = append(establishedOn, hit)
|
||||
}
|
||||
if bestState == "" || sessionStateRank(sess.State) > sessionStateRank(bestState) {
|
||||
bestState = sess.State
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if len(establishedOn) > 0 {
|
||||
bestState = "Established"
|
||||
labels := make([]string, 0, len(establishedOn))
|
||||
for _, e := range establishedOn {
|
||||
labels = append(labels, e.Label)
|
||||
}
|
||||
connectedLabel = strings.Join(labels, ", ")
|
||||
if len(establishedOn) == 1 {
|
||||
connectedID = establishedOn[0].SpeakerID
|
||||
}
|
||||
} else if len(on) == 1 {
|
||||
connectedID = on[0].SpeakerID
|
||||
connectedLabel = on[0].Label
|
||||
}
|
||||
|
||||
if peer.SpeakerID != nil && strings.TrimSpace(*peer.SpeakerID) != "" && len(establishedOn) > 0 {
|
||||
want := strings.TrimSpace(*peer.SpeakerID)
|
||||
found := false
|
||||
for _, e := range establishedOn {
|
||||
if strings.EqualFold(strings.TrimSpace(e.SpeakerID), want) {
|
||||
found = true
|
||||
break
|
||||
}
|
||||
}
|
||||
mismatch = !found
|
||||
}
|
||||
return bestState, connectedID, connectedLabel, establishedOn, on, mismatch
|
||||
}
|
||||
|
||||
func peerSessionMatches(sess birdfmt.BGPSession, protoName string, neighbor netip.Addr, hasNeighbor bool) bool {
|
||||
if strings.EqualFold(strings.TrimSpace(sess.Name), protoName) {
|
||||
return true
|
||||
}
|
||||
if !hasNeighbor || strings.TrimSpace(sess.Neighbor) == "" {
|
||||
return false
|
||||
}
|
||||
addr, err := netip.ParseAddr(strings.TrimSpace(sess.Neighbor))
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
return addr == neighbor
|
||||
}
|
||||
|
||||
func sessionStateRank(state string) int {
|
||||
switch strings.ToLower(strings.TrimSpace(state)) {
|
||||
case "established":
|
||||
return 100
|
||||
case "openconfirm", "opensent":
|
||||
return 80
|
||||
case "active", "connect":
|
||||
return 60
|
||||
case "idle":
|
||||
return 20
|
||||
default:
|
||||
return 10
|
||||
}
|
||||
}
|
||||
|
||||
func applyPeerLiveFields(row map[string]any, peer *store.BGPPeer, views []speakerBGPLive) {
|
||||
state, connID, connLabel, establishedOn, on, mismatch := matchPeerOnSpeakers(peer, views)
|
||||
if len(on) > 0 {
|
||||
if state != "" {
|
||||
row["session_state"] = state
|
||||
}
|
||||
row["connected_speaker_id"] = peerLiveSpeakerIDOrNull(connID)
|
||||
row["connected_speaker_label"] = connLabel
|
||||
row["session_on_speakers"] = on
|
||||
if len(establishedOn) > 0 {
|
||||
row["established_on_speakers"] = establishedOn
|
||||
} else {
|
||||
row["established_on_speakers"] = []peerSessionOnSpeaker{}
|
||||
}
|
||||
row["session_conflict"] = false
|
||||
row["session_mismatch"] = mismatch
|
||||
return
|
||||
}
|
||||
row["session_on_speakers"] = []peerSessionOnSpeaker{}
|
||||
row["established_on_speakers"] = []peerSessionOnSpeaker{}
|
||||
row["session_conflict"] = false
|
||||
row["session_mismatch"] = false
|
||||
}
|
||||
|
||||
func peerLiveSpeakerIDOrNull(id string) any {
|
||||
if strings.TrimSpace(id) == "" {
|
||||
return nil
|
||||
}
|
||||
return id
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
package httpapi
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"evobgp/internal/birdfmt"
|
||||
"evobgp/internal/store"
|
||||
)
|
||||
|
||||
func TestMatchPeerOnSpeakers_establishedOnReplica(t *testing.T) {
|
||||
peer := &store.BGPPeer{
|
||||
ID: "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee",
|
||||
Neighbor: "198.51.100.2",
|
||||
}
|
||||
views := []speakerBGPLive{
|
||||
{
|
||||
SpeakerID: "master-id",
|
||||
Label: "CP · bgp.shz.su",
|
||||
Sessions: []birdfmt.BGPSession{
|
||||
{Name: birdfmt.PeerProtocolName(peer.ID), Neighbor: "198.51.100.2", State: "Idle"},
|
||||
},
|
||||
},
|
||||
{
|
||||
SpeakerID: "replica-id",
|
||||
Label: "bgp2.shz.su",
|
||||
Sessions: []birdfmt.BGPSession{
|
||||
{Name: birdfmt.PeerProtocolName(peer.ID), Neighbor: "198.51.100.2", State: "Established"},
|
||||
},
|
||||
},
|
||||
}
|
||||
state, connID, connLabel, established, on, mismatch := matchPeerOnSpeakers(peer, views)
|
||||
if state != "Established" || connID != "replica-id" || connLabel != "bgp2.shz.su" {
|
||||
t.Fatalf("got state=%q conn=%q label=%q", state, connID, connLabel)
|
||||
}
|
||||
if mismatch || len(on) != 2 || len(established) != 1 {
|
||||
t.Fatalf("on=%+v established=%+v mismatch=%v", on, established, mismatch)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMatchPeerOnSpeakers_multipleEstablished(t *testing.T) {
|
||||
peer := &store.BGPPeer{ID: "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee", Neighbor: "198.51.100.2"}
|
||||
views := []speakerBGPLive{
|
||||
{SpeakerID: "a", Label: "n1", Sessions: []birdfmt.BGPSession{{Name: birdfmt.PeerProtocolName(peer.ID), State: "Established"}}},
|
||||
{SpeakerID: "b", Label: "n2", Sessions: []birdfmt.BGPSession{{Name: birdfmt.PeerProtocolName(peer.ID), State: "Established"}}},
|
||||
}
|
||||
_, connID, label, established, _, mismatch := matchPeerOnSpeakers(peer, views)
|
||||
if mismatch || connID != "" || label != "n1, n2" || len(established) != 2 {
|
||||
t.Fatalf("connID=%q label=%q established=%+v mismatch=%v", connID, label, established, mismatch)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMatchPeerOnSpeakers_mismatchConfiguredSpeaker(t *testing.T) {
|
||||
replica := "replica-id"
|
||||
peer := &store.BGPPeer{
|
||||
ID: "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee",
|
||||
Neighbor: "198.51.100.2",
|
||||
SpeakerID: &replica,
|
||||
}
|
||||
views := []speakerBGPLive{
|
||||
{SpeakerID: "master-id", Label: "CP", Sessions: []birdfmt.BGPSession{{Name: birdfmt.PeerProtocolName(peer.ID), State: "Established"}}},
|
||||
{SpeakerID: replica, Label: "bgp2", Sessions: []birdfmt.BGPSession{{Name: birdfmt.PeerProtocolName(peer.ID), State: "Idle"}}},
|
||||
}
|
||||
_, _, _, _, _, mismatch := matchPeerOnSpeakers(peer, views)
|
||||
if !mismatch {
|
||||
t.Fatal("expected mismatch when configured replica has no Established")
|
||||
}
|
||||
}
|
||||
@@ -283,13 +283,14 @@ func (s *Server) handleListPeers(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
allPeers := s.store.ListPeers(a.TenantID)
|
||||
page, next, more := store.PaginateOffset(allPeers, r.URL.Query().Get("cursor"), parseListLimit(r))
|
||||
liveStates := s.liveBGPProtocolStates(r)
|
||||
fresh := r != nil && strings.EqualFold(strings.TrimSpace(r.URL.Query().Get("live")), "1")
|
||||
ctx, cancel := context.WithTimeout(r.Context(), 12*time.Second)
|
||||
defer cancel()
|
||||
liveViews := s.collectSpeakerBGPLive(ctx, a.TenantID, fresh)
|
||||
items := make([]map[string]any, 0, len(page))
|
||||
for _, p := range page {
|
||||
row := peerJSON(p)
|
||||
if st, ok := liveStates[peerProtocolNameForID(p.ID)]; ok && strings.TrimSpace(st) != "" {
|
||||
row["session_state"] = strings.TrimSpace(st)
|
||||
}
|
||||
applyPeerLiveFields(row, p, liveViews)
|
||||
items = append(items, row)
|
||||
}
|
||||
writeJSON(w, http.StatusOK, map[string]any{
|
||||
@@ -366,16 +367,9 @@ func extractBGPSessionState(line string) string {
|
||||
return ""
|
||||
}
|
||||
|
||||
// peerProtocolNameForID must stay in sync with pipeline peer protocol naming.
|
||||
// peerProtocolNameForID forwards to birdfmt for tests and legacy callers.
|
||||
func peerProtocolNameForID(peerID string) string {
|
||||
s := strings.ReplaceAll(strings.TrimSpace(peerID), "-", "")
|
||||
if len(s) > 16 {
|
||||
s = s[:16]
|
||||
}
|
||||
if s == "" {
|
||||
s = "x"
|
||||
}
|
||||
return "evobgp_p_" + s
|
||||
return birdfmt.PeerProtocolName(peerID)
|
||||
}
|
||||
|
||||
func (s *Server) handleListSpeakers(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
Reference in New Issue
Block a user