Files

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)