package mesh import ( "context" "fmt" "strings" "sync" ) type FabricSessionPeerManager struct { mu sync.Mutex sessions map[string]*FabricSessionPump } type FabricSessionPeerTarget struct { PeerID string BaseURL string Options FabricSessionDialOptions Pump FabricSessionPumpOptions } func NewFabricSessionPeerManager() *FabricSessionPeerManager { return &FabricSessionPeerManager{ sessions: map[string]*FabricSessionPump{}, } } func (m *FabricSessionPeerManager) Get(ctx context.Context, target FabricSessionPeerTarget) (*FabricSessionPump, error) { if m == nil { return nil, fmt.Errorf("fabric session peer manager is nil") } key, err := fabricSessionPeerKey(target) if err != nil { return nil, err } m.mu.Lock() if pump := m.sessions[key]; pump != nil { if pump.Closed() { delete(m.sessions, key) } else { m.mu.Unlock() return pump, nil } } m.mu.Unlock() session, _, err := NewClient(target.BaseURL).OpenFabricSession(ctx, target.Options) if err != nil { return nil, err } pump := session.StartPump(ctx, target.Pump) m.mu.Lock() if existing := m.sessions[key]; existing != nil { if existing.Closed() { delete(m.sessions, key) } else { m.mu.Unlock() _ = pump.Close() return existing, nil } } if m.sessions == nil { m.sessions = map[string]*FabricSessionPump{} } m.sessions[key] = pump m.mu.Unlock() return pump, nil } func (m *FabricSessionPeerManager) ClosePeer(target FabricSessionPeerTarget) error { if m == nil { return nil } key, err := fabricSessionPeerKey(target) if err != nil { return err } m.mu.Lock() pump := m.sessions[key] delete(m.sessions, key) m.mu.Unlock() if pump == nil { return nil } return pump.Close() } func (m *FabricSessionPeerManager) Close() error { if m == nil { return nil } m.mu.Lock() sessions := m.sessions m.sessions = map[string]*FabricSessionPump{} m.mu.Unlock() var firstErr error for _, pump := range sessions { if err := pump.Close(); err != nil && firstErr == nil { firstErr = err } } return firstErr } func fabricSessionPeerKey(target FabricSessionPeerTarget) (string, error) { peerID := strings.TrimSpace(target.PeerID) baseURL := strings.TrimRight(strings.TrimSpace(target.BaseURL), "/") if peerID == "" || baseURL == "" { return "", fmt.Errorf("fabric session peer id and base url are required") } return peerID + "\x00" + baseURL, nil }