pool.mx raw

   1  package relay
   2  
   3  // Pool manages connections to multiple relays.
   4  type Pool struct {
   5  	conns map[string]*Conn
   6  }
   7  
   8  // NewPool creates a relay pool.
   9  func NewPool() (p *Pool) {
  10  	return &Pool{
  11  		conns: map[string]*Conn{},
  12  	}
  13  }
  14  
  15  // Connect gets or creates a connection to a relay.
  16  //
  17  // An existing Conn is reused even when it is currently closed: Conn.maybeReconnect
  18  // redials the SAME object, and its flushSubs re-establishes every subscription
  19  // held in c.subs. Replacing it instead orphaned those subscriptions - the new
  20  // Conn started with an empty registry, nothing re-subscribed the entries the
  21  // consumer still had registered, and that relay went silently dead for them.
  22  func (p *Pool) Connect(url string) (c *Conn) {
  23  	if existing, ok := p.conns[url]; ok {
  24  		return existing
  25  	}
  26  	c = Dial(url)
  27  	p.conns[url] = c
  28  	return c
  29  }
  30  
  31  // Get returns an existing connection, or nil.
  32  func (p *Pool) Get(url string) (c *Conn) {
  33  	got, ok := p.conns[url]
  34  	if !ok || !got.IsOpen() {
  35  		return nil
  36  	}
  37  	return got
  38  }
  39  
  40  // Disconnect closes and removes a connection.
  41  func (p *Pool) Disconnect(url string) {
  42  	if c, ok2 := p.conns[url]; ok2 {
  43  		c.Close()
  44  		delete(p.conns, url)
  45  	}
  46  }
  47  
  48  // CloseAll closes all connections.
  49  func (p *Pool) CloseAll() {
  50  	for url, c := range p.conns {
  51  		c.Close()
  52  		delete(p.conns, url)
  53  	}
  54  }
  55  
  56  // URLs returns all connected relay URLs.
  57  func (p *Pool) URLs() (ss []string) {
  58  	var out []string
  59  	for url, c := range p.conns {
  60  		if c.IsOpen() {
  61  			out = push(out, url)
  62  		}
  63  	}
  64  	return out
  65  }
  66  
  67  // AllConns returns all connections regardless of state.
  68  func (p *Pool) AllConns() (cs []*Conn) {
  69  	p.evictClosed()
  70  	var out []*Conn
  71  	for _, c := range p.conns {
  72  		out = push(out, c)
  73  	}
  74  	return out
  75  }
  76  
  77  // evictClosed removes closed connections that hold no subscriptions. A closed
  78  // Conn that still holds them is mid-reconnect: it redials itself and flushes
  79  // them again, so dropping it here would lose them for good.
  80  func (p *Pool) evictClosed() {
  81  	for url, c := range p.conns {
  82  		if c.state == StateClosed && len(c.subs) == 0 {
  83  			delete(p.conns, url)
  84  		}
  85  	}
  86  }
  87