Files

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)
}
}