2020-05-01 05:47:46 +00:00
|
|
|
package api
|
2017-02-02 03:12:13 +00:00
|
|
|
|
2017-03-22 21:43:36 +00:00
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"fmt"
|
2020-05-01 05:47:46 +00:00
|
|
|
"log"
|
2017-03-22 21:43:36 +00:00
|
|
|
"time"
|
|
|
|
)
|
2017-02-12 20:39:50 +00:00
|
|
|
|
|
|
|
const (
|
|
|
|
initialDomains = 0
|
|
|
|
incrementDomains = 0
|
|
|
|
)
|
|
|
|
|
2017-02-12 04:13:29 +00:00
|
|
|
//Table maintains the set of connections
|
|
|
|
type Table struct {
|
2017-02-12 20:39:50 +00:00
|
|
|
connections map[*Connection][]string
|
2020-05-01 05:47:46 +00:00
|
|
|
Domains map[string]*DomainLoadBalance
|
2017-02-19 20:05:06 +00:00
|
|
|
register chan *Registration
|
2017-02-12 20:39:50 +00:00
|
|
|
unregister chan *Connection
|
|
|
|
domainAnnounce chan *DomainMapping
|
|
|
|
domainRevoke chan *DomainMapping
|
2017-03-11 20:28:49 +00:00
|
|
|
dwell int
|
|
|
|
idle int
|
2020-04-30 10:43:36 +00:00
|
|
|
balanceMethod string
|
2017-02-02 03:12:13 +00:00
|
|
|
}
|
|
|
|
|
2017-02-12 04:13:29 +00:00
|
|
|
//NewTable -- consructor
|
2020-04-30 10:43:36 +00:00
|
|
|
func NewTable(dwell, idle int, balanceMethod string) (p *Table) {
|
2017-02-12 20:39:50 +00:00
|
|
|
p = new(Table)
|
|
|
|
p.connections = make(map[*Connection][]string)
|
2020-05-01 05:47:46 +00:00
|
|
|
p.Domains = make(map[string]*DomainLoadBalance)
|
2017-02-19 20:05:06 +00:00
|
|
|
p.register = make(chan *Registration)
|
2017-02-12 20:39:50 +00:00
|
|
|
p.unregister = make(chan *Connection)
|
|
|
|
p.domainAnnounce = make(chan *DomainMapping)
|
|
|
|
p.domainRevoke = make(chan *DomainMapping)
|
2017-03-11 20:28:49 +00:00
|
|
|
p.dwell = dwell
|
|
|
|
p.idle = idle
|
2020-04-30 10:43:36 +00:00
|
|
|
p.balanceMethod = balanceMethod
|
2017-02-12 20:39:50 +00:00
|
|
|
return
|
2017-02-02 03:12:13 +00:00
|
|
|
}
|
|
|
|
|
2017-02-15 23:53:34 +00:00
|
|
|
//Connections Property
|
2017-03-22 22:33:09 +00:00
|
|
|
func (c *Table) Connections() map[*Connection][]string {
|
|
|
|
return c.connections
|
2017-02-15 23:53:34 +00:00
|
|
|
}
|
|
|
|
|
2017-03-25 20:01:57 +00:00
|
|
|
//ConnByDomain -- Obtains a connection from a domain announcement. A domain may be announced more than once
|
|
|
|
//if that is the case the system stores these connections and then sends traffic back round-robin
|
|
|
|
//back to the WSS connections
|
2017-03-22 22:33:09 +00:00
|
|
|
func (c *Table) ConnByDomain(domain string) (*Connection, bool) {
|
2020-05-01 05:47:46 +00:00
|
|
|
for dn := range c.Domains {
|
|
|
|
log.Println(dn, domain)
|
2017-03-25 20:01:57 +00:00
|
|
|
}
|
2020-05-01 05:47:46 +00:00
|
|
|
if domainsLB, ok := c.Domains[domain]; ok {
|
|
|
|
log.Println("found")
|
2017-03-25 20:01:57 +00:00
|
|
|
conn := domainsLB.NextMember()
|
|
|
|
return conn, ok
|
|
|
|
}
|
|
|
|
return nil, false
|
2017-02-14 02:36:01 +00:00
|
|
|
}
|
|
|
|
|
2017-02-19 20:32:03 +00:00
|
|
|
//reaper --
|
|
|
|
func (c *Table) reaper(delay int, idle int) {
|
2017-02-19 21:51:54 +00:00
|
|
|
_ = "breakpoint"
|
2017-02-19 20:32:03 +00:00
|
|
|
for {
|
2020-05-01 05:47:46 +00:00
|
|
|
log.Println("Reaper waiting for ", delay, " seconds")
|
2017-02-19 20:32:03 +00:00
|
|
|
time.Sleep(time.Duration(delay) * time.Second)
|
|
|
|
|
2020-05-01 05:47:46 +00:00
|
|
|
log.Println("Running scanning ", len(c.connections))
|
2017-02-19 20:32:03 +00:00
|
|
|
for d := range c.connections {
|
2017-03-22 23:45:47 +00:00
|
|
|
if !d.State() {
|
2017-02-19 20:32:03 +00:00
|
|
|
if time.Since(d.lastUpdate).Seconds() > float64(idle) {
|
2020-05-01 05:47:46 +00:00
|
|
|
log.Println("reaper removing ", d.lastUpdate, time.Since(d.lastUpdate).Seconds())
|
2017-02-19 20:32:03 +00:00
|
|
|
delete(c.connections, d)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2017-03-13 21:46:11 +00:00
|
|
|
//GetConnection -- find connection by server-id
|
2017-03-22 22:33:09 +00:00
|
|
|
func (c *Table) GetConnection(serverID int64) (*Connection, error) {
|
2017-03-13 21:46:11 +00:00
|
|
|
for conn := range c.connections {
|
|
|
|
if conn.ConnectionID() == serverID {
|
2017-03-22 22:33:09 +00:00
|
|
|
return conn, nil
|
2017-03-13 21:46:11 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2017-03-22 22:33:09 +00:00
|
|
|
return nil, fmt.Errorf("Server-id %d not found", serverID)
|
2017-03-13 21:46:11 +00:00
|
|
|
}
|
|
|
|
|
2017-02-12 04:13:29 +00:00
|
|
|
//Run -- Execute
|
2020-04-30 10:43:36 +00:00
|
|
|
func (c *Table) Run(ctx context.Context) {
|
2020-05-01 05:47:46 +00:00
|
|
|
log.Println("ConnectionTable starting")
|
2017-02-19 20:32:03 +00:00
|
|
|
|
2017-03-11 20:28:49 +00:00
|
|
|
go c.reaper(c.dwell, c.idle)
|
2017-02-19 20:32:03 +00:00
|
|
|
|
2017-02-02 03:12:13 +00:00
|
|
|
for {
|
|
|
|
select {
|
2017-02-26 05:17:39 +00:00
|
|
|
|
|
|
|
case <-ctx.Done():
|
2020-05-01 05:47:46 +00:00
|
|
|
log.Println("Cancel signal hit")
|
2017-02-26 05:17:39 +00:00
|
|
|
return
|
|
|
|
|
2017-02-19 20:05:06 +00:00
|
|
|
case registration := <-c.register:
|
2020-05-01 05:47:46 +00:00
|
|
|
log.Println("register fired")
|
2017-02-19 20:05:06 +00:00
|
|
|
|
2017-03-25 22:22:28 +00:00
|
|
|
connection := NewConnection(c, registration.conn, registration.source, registration.initialDomains,
|
|
|
|
registration.connectionTrack, registration.serverName)
|
2017-02-13 21:59:10 +00:00
|
|
|
c.connections[connection] = make([]string, initialDomains)
|
2017-02-19 20:05:06 +00:00
|
|
|
registration.commCh <- true
|
2017-02-02 03:12:13 +00:00
|
|
|
|
2017-02-12 20:39:50 +00:00
|
|
|
// handle initial domain additions
|
|
|
|
for _, domain := range connection.initialDomains {
|
|
|
|
// add to the domains regirstation
|
|
|
|
|
2020-04-30 11:11:03 +00:00
|
|
|
newDomain := domain
|
2020-05-01 05:47:46 +00:00
|
|
|
log.Println("adding domain ", newDomain, " to connection ", connection.conn.RemoteAddr().String())
|
2017-03-25 20:01:57 +00:00
|
|
|
|
|
|
|
//check to see if domain is already present.
|
2020-05-01 05:47:46 +00:00
|
|
|
if _, ok := c.Domains[newDomain]; ok {
|
2017-03-25 20:01:57 +00:00
|
|
|
|
|
|
|
//append to a list of connections for that domain
|
2020-05-01 05:47:46 +00:00
|
|
|
c.Domains[newDomain].AddConnection(connection)
|
2017-03-25 20:01:57 +00:00
|
|
|
} else {
|
|
|
|
//if not, then add as the 1st to the list of connections
|
2020-05-01 05:47:46 +00:00
|
|
|
c.Domains[newDomain] = NewDomainLoadBalance(c.balanceMethod)
|
|
|
|
c.Domains[newDomain].AddConnection(connection)
|
2017-03-25 20:01:57 +00:00
|
|
|
}
|
2017-02-12 20:39:50 +00:00
|
|
|
|
|
|
|
// add to the connection domain list
|
|
|
|
s := c.connections[connection]
|
|
|
|
c.connections[connection] = append(s, newDomain)
|
2017-02-02 03:12:13 +00:00
|
|
|
}
|
2017-02-19 20:05:06 +00:00
|
|
|
go connection.Writer()
|
2017-03-03 00:47:59 +00:00
|
|
|
go connection.Reader(ctx)
|
2017-02-12 20:39:50 +00:00
|
|
|
|
2017-02-02 03:12:13 +00:00
|
|
|
case connection := <-c.unregister:
|
2020-05-01 05:47:46 +00:00
|
|
|
log.Println("closing connection ", connection.conn.RemoteAddr().String())
|
2017-03-25 20:01:57 +00:00
|
|
|
|
|
|
|
//does connection exist in the connection table -- should never be an issue
|
2017-02-02 03:12:13 +00:00
|
|
|
if _, ok := c.connections[connection]; ok {
|
2017-03-25 20:01:57 +00:00
|
|
|
|
|
|
|
//iterate over the connections for the domain
|
2017-02-12 20:39:50 +00:00
|
|
|
for _, domain := range c.connections[connection] {
|
2020-05-01 05:47:46 +00:00
|
|
|
log.Println("remove domain", domain)
|
2017-03-25 20:01:57 +00:00
|
|
|
|
|
|
|
//removing domain, make sure it is present (should never be a problem)
|
2020-05-01 05:47:46 +00:00
|
|
|
if _, ok := c.Domains[domain]; ok {
|
2017-03-25 20:01:57 +00:00
|
|
|
|
2020-05-01 05:47:46 +00:00
|
|
|
domainLB := c.Domains[domain]
|
2017-03-25 20:01:57 +00:00
|
|
|
domainLB.RemoveConnection(connection)
|
|
|
|
|
|
|
|
//check to see if domain is free of connections, if yes, delete map entry
|
|
|
|
if domainLB.count > 0 {
|
|
|
|
//ignore...perhaps we will do something here dealing wtih the lb method
|
|
|
|
} else {
|
2020-05-01 05:47:46 +00:00
|
|
|
delete(c.Domains, domain)
|
2017-03-25 20:01:57 +00:00
|
|
|
}
|
2017-02-12 20:39:50 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2017-02-19 20:05:06 +00:00
|
|
|
//delete(c.connections, connection)
|
|
|
|
//close(connection.send)
|
2017-02-02 03:12:13 +00:00
|
|
|
}
|
2017-02-12 20:39:50 +00:00
|
|
|
|
|
|
|
case domainMapping := <-c.domainAnnounce:
|
2020-05-01 05:47:46 +00:00
|
|
|
log.Println("domainMapping fired ", domainMapping)
|
2017-02-12 20:39:50 +00:00
|
|
|
//check to make sure connection is already regiered, you can no register a domain without an apporved connection
|
|
|
|
//if connection, ok := connections[domainMapping.connection]; ok {
|
|
|
|
|
|
|
|
//} else {
|
|
|
|
|
|
|
|
//}
|
|
|
|
|
2017-02-02 03:12:13 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2017-02-12 04:13:29 +00:00
|
|
|
|
|
|
|
//Register -- Property
|
2017-03-22 22:33:09 +00:00
|
|
|
func (c *Table) Register() chan *Registration {
|
|
|
|
return c.register
|
2017-02-12 04:13:29 +00:00
|
|
|
}
|