135 lines
2.7 KiB
Go
135 lines
2.7 KiB
Go
package xmpp
|
|
|
|
import (
|
|
"fmt"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"git.i3omb.com/gronod/ha-n95-local-control/internal/session"
|
|
)
|
|
|
|
const (
|
|
// PingPeriod is the interval between controller announce/keepalive pings.
|
|
PingPeriod = 60 * time.Second
|
|
// PingResultTimeout is how long to wait for a matching ping result before
|
|
// treating the session as half-open.
|
|
PingResultTimeout = 12 * time.Second
|
|
)
|
|
|
|
var pingCounter atomic.Uint64
|
|
|
|
// pingLoop sends periodic controller pings and waits for matching results.
|
|
type pingLoop struct {
|
|
c *Conn
|
|
mu sync.Mutex
|
|
pending string // current outstanding ping id
|
|
announced bool // AnnounceOK already emitted for this connection
|
|
resultCh chan string
|
|
stopCh chan struct{}
|
|
doneCh chan struct{}
|
|
}
|
|
|
|
func newPingLoop(c *Conn) *pingLoop {
|
|
return &pingLoop{
|
|
c: c,
|
|
resultCh: make(chan string, 1),
|
|
stopCh: make(chan struct{}),
|
|
doneCh: make(chan struct{}),
|
|
}
|
|
}
|
|
|
|
func (p *pingLoop) nextID() string {
|
|
return fmt.Sprintf("%d", pingCounter.Add(1))
|
|
}
|
|
|
|
func (p *pingLoop) run() {
|
|
defer close(p.doneCh)
|
|
send := func(b []byte) error { return p.c.registry.Send(p.c.fullJID, p.c.gen, b) }
|
|
first := true
|
|
for {
|
|
if !first {
|
|
select {
|
|
case <-p.stopCh:
|
|
return
|
|
case <-p.c.clock.After(PingPeriod):
|
|
}
|
|
}
|
|
first = false
|
|
|
|
id := p.nextID()
|
|
p.mu.Lock()
|
|
p.pending = id
|
|
p.mu.Unlock()
|
|
|
|
ping := fmt.Sprintf(`<iq id="%s" to="%s" from="%s" type="get"><ping xmlns="urn:xmpp:ping"/></iq>`,
|
|
id, escapeXMLAttr(p.c.fullJID), escapeXMLAttr(p.c.cfg.ControllerJID))
|
|
if err := send([]byte(ping)); err != nil {
|
|
return
|
|
}
|
|
|
|
timeout := p.c.clock.After(PingResultTimeout)
|
|
waitResult:
|
|
for {
|
|
select {
|
|
case <-p.stopCh:
|
|
return
|
|
case <-timeout:
|
|
p.mu.Lock()
|
|
still := p.pending == id
|
|
p.mu.Unlock()
|
|
if still {
|
|
p.c.endSession(session.ReasonPingTimeout)
|
|
}
|
|
return
|
|
case rid := <-p.resultCh:
|
|
p.mu.Lock()
|
|
if p.pending != rid {
|
|
p.mu.Unlock()
|
|
continue
|
|
}
|
|
p.pending = ""
|
|
announced := p.announced
|
|
p.announced = true
|
|
p.mu.Unlock()
|
|
if !announced {
|
|
p.c.bus.AnnounceOK(session.ReadyEvent{
|
|
Generation: p.c.gen,
|
|
JID: p.c.fullJID,
|
|
Serial: p.c.serial,
|
|
Send: send,
|
|
})
|
|
}
|
|
break waitResult
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// result reports a matching ping result id. It returns true when the id is the
|
|
// currently outstanding ping.
|
|
func (p *pingLoop) result(id string) bool {
|
|
p.mu.Lock()
|
|
pending := p.pending == id
|
|
p.mu.Unlock()
|
|
if !pending {
|
|
return false
|
|
}
|
|
select {
|
|
case p.resultCh <- id:
|
|
return true
|
|
case <-p.doneCh:
|
|
return false
|
|
}
|
|
}
|
|
|
|
// stop signals the ping loop to exit.
|
|
func (p *pingLoop) stop() {
|
|
select {
|
|
case <-p.stopCh:
|
|
return
|
|
default:
|
|
close(p.stopCh)
|
|
}
|
|
}
|