224 lines
5.2 KiB
Go
224 lines
5.2 KiB
Go
// Package session owns the XMPP generation registry, observer bus, and the
|
|
// diagnostic sink type used across the XMPP and robot layers.
|
|
package session
|
|
|
|
import (
|
|
"errors"
|
|
"sync"
|
|
)
|
|
|
|
// ErrStale is returned by Registry.Send when the requested generation is no
|
|
// longer current for the JID.
|
|
var ErrStale = errors.New("stale generation")
|
|
|
|
// Reasons carried by DownEvent.
|
|
const (
|
|
ReasonTCPClose = "tcp-close"
|
|
ReasonStreamClose = "stream-close"
|
|
ReasonPingTimeout = "ping-timeout"
|
|
ReasonMalformedXML = "malformed-xml"
|
|
ReasonReplaced = "replaced"
|
|
)
|
|
|
|
// Diagnostic directions.
|
|
const (
|
|
DirectionIn = "in"
|
|
DirectionOut = "out"
|
|
)
|
|
|
|
// Registry holds at most one current generation per full bot JID. Bind
|
|
// atomically replaces any existing slot; the old slot is closed and the
|
|
// previous generation becomes stale.
|
|
type Registry struct {
|
|
mu sync.Mutex
|
|
current map[string]*slot
|
|
}
|
|
|
|
type slot struct {
|
|
gen uint64
|
|
send func([]byte) error
|
|
close func()
|
|
}
|
|
|
|
// NewRegistry creates an empty registry.
|
|
func NewRegistry() *Registry {
|
|
return &Registry{current: make(map[string]*slot)}
|
|
}
|
|
|
|
// Bind installs a new generation for jid. It returns the new generation and
|
|
// whether an existing slot was replaced. The old slot is closed after the
|
|
// registry has been updated.
|
|
func (r *Registry) Bind(jid string, send func([]byte) error, close func()) (gen uint64, replaced bool) {
|
|
r.mu.Lock()
|
|
old := r.current[jid]
|
|
var next uint64
|
|
if old != nil {
|
|
replaced = true
|
|
next = old.gen + 1
|
|
} else {
|
|
next = 1
|
|
}
|
|
r.current[jid] = &slot{gen: next, send: send, close: close}
|
|
r.mu.Unlock()
|
|
|
|
if old != nil && old.close != nil {
|
|
old.close()
|
|
}
|
|
return next, replaced
|
|
}
|
|
|
|
// Current returns the current generation for jid, if any.
|
|
func (r *Registry) Current(jid string) (gen uint64, ok bool) {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
s := r.current[jid]
|
|
if s == nil {
|
|
return 0, false
|
|
}
|
|
return s.gen, true
|
|
}
|
|
|
|
// Send delivers xml to the current generation. It returns ErrStale if the
|
|
// generation is no longer current. The lock is released before calling the
|
|
// slot's send function so a slow or blocked write cannot stall the registry.
|
|
func (r *Registry) Send(jid string, gen uint64, xml []byte) error {
|
|
r.mu.Lock()
|
|
s := r.current[jid]
|
|
r.mu.Unlock()
|
|
if s == nil || s.gen != gen {
|
|
return ErrStale
|
|
}
|
|
return s.send(xml)
|
|
}
|
|
|
|
// Remove deletes the slot for jid only if it matches gen. It returns true when
|
|
// the slot existed and matched, which a read loop uses after emitting its own
|
|
// SessionDown for the current generation.
|
|
func (r *Registry) Remove(jid string, gen uint64) bool {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
s := r.current[jid]
|
|
if s != nil && s.gen == gen {
|
|
delete(r.current, jid)
|
|
return true
|
|
}
|
|
return false
|
|
}
|
|
|
|
// CloseAll calls every current slot's close function without removing it. The
|
|
// read loops observe the closed socket and emit SessionDown. New slots cannot
|
|
// be bound while CloseAll runs because the caller owns shutdown order.
|
|
func (r *Registry) CloseAll() []func() {
|
|
r.mu.Lock()
|
|
closers := make([]func(), 0, len(r.current))
|
|
for _, s := range r.current {
|
|
closers = append(closers, s.close)
|
|
}
|
|
r.mu.Unlock()
|
|
for _, c := range closers {
|
|
c()
|
|
}
|
|
return closers
|
|
}
|
|
|
|
// ReadyEvent is emitted when a full handshake reaches READY.
|
|
type ReadyEvent struct {
|
|
Generation uint64
|
|
JID string
|
|
Serial string
|
|
Send func([]byte) error
|
|
}
|
|
|
|
// DownEvent is emitted when a generation ends.
|
|
type DownEvent struct {
|
|
Generation uint64
|
|
JID string
|
|
Serial string
|
|
Reason string
|
|
}
|
|
|
|
// StanzaEvent is emitted for every non-ping post-READY stanza.
|
|
type StanzaEvent struct {
|
|
Generation uint64
|
|
JID string
|
|
Serial string
|
|
Stanza []byte
|
|
}
|
|
|
|
// Observer receives session lifecycle events. Implementations may implement
|
|
// only the methods they care about by providing no-op bodies.
|
|
type Observer interface {
|
|
SessionReady(e ReadyEvent)
|
|
AnnounceOK(e ReadyEvent)
|
|
SessionDown(e DownEvent)
|
|
Stanza(e StanzaEvent)
|
|
}
|
|
|
|
// Bus fans observer calls out to registered observers.
|
|
type Bus struct {
|
|
mu sync.RWMutex
|
|
obs []Observer
|
|
}
|
|
|
|
// NewBus creates an empty bus.
|
|
func NewBus() *Bus {
|
|
return &Bus{}
|
|
}
|
|
|
|
// Register adds an observer. Observers are called in registration order.
|
|
func (b *Bus) Register(o Observer) {
|
|
b.mu.Lock()
|
|
defer b.mu.Unlock()
|
|
b.obs = append(b.obs, o)
|
|
}
|
|
|
|
func (b *Bus) SessionReady(e ReadyEvent) {
|
|
b.mu.RLock()
|
|
obs := append([]Observer(nil), b.obs...)
|
|
b.mu.RUnlock()
|
|
for _, o := range obs {
|
|
o.SessionReady(e)
|
|
}
|
|
}
|
|
|
|
func (b *Bus) AnnounceOK(e ReadyEvent) {
|
|
b.mu.RLock()
|
|
obs := append([]Observer(nil), b.obs...)
|
|
b.mu.RUnlock()
|
|
for _, o := range obs {
|
|
o.AnnounceOK(e)
|
|
}
|
|
}
|
|
|
|
func (b *Bus) SessionDown(e DownEvent) {
|
|
b.mu.RLock()
|
|
obs := append([]Observer(nil), b.obs...)
|
|
b.mu.RUnlock()
|
|
for _, o := range obs {
|
|
o.SessionDown(e)
|
|
}
|
|
}
|
|
|
|
func (b *Bus) Stanza(e StanzaEvent) {
|
|
b.mu.RLock()
|
|
obs := append([]Observer(nil), b.obs...)
|
|
b.mu.RUnlock()
|
|
for _, o := range obs {
|
|
o.Stanza(e)
|
|
}
|
|
}
|
|
|
|
// Diagnostic carries a redacted XMPP fragment for the raw diagnostic topic.
|
|
// Auth material never reaches a Diagnostic.
|
|
type Diagnostic struct {
|
|
Serial string
|
|
Direction string
|
|
Kind string
|
|
Generation uint64
|
|
XML []byte
|
|
Reason string
|
|
}
|
|
|
|
// DiagnosticSink receives diagnostics. A nil sink drops the event.
|
|
type DiagnosticSink func(Diagnostic)
|