package api import "sync" // Hub wakes live listeners when a tenant's data changes. // // It is a DOORBELL, not a delivery service: a notification carries a client id // and nothing else, and every listener answers it by running the same keyset // query a polling client would. That is the whole design decision, and it buys // three things that a hub carrying the rows does not: // // - One query path. The stream and the poll cannot disagree about what an // arrival looks like, because there is only one piece of code that reads // one. // - No lost events. A subscriber that is mid-reconnect, or slow, or was not // listening yet, misses a doorbell and loses nothing - its next query // starts from its own cursor and picks up everything in between. A hub // that pushed rows would have to buffer per subscriber and decide what to // drop, which is a queue, and we already have a durable one. // - Degrades to polling. If a second server instance is ever added, its // ingest rings a doorbell this process never hears. The stream keeps a // slow fallback tick for exactly that, so the failure mode is latency, // not silence. // // Notify is called from the MQTT consumer's goroutine and must never block it: // a slow subscriber must not be able to stall ingest for the whole estate. type Hub struct { mu sync.Mutex next int subs map[int]chan string } func NewHub() *Hub { return &Hub{subs: map[int]chan string{}} } // Subscribe returns a channel of client ids and a function to release it. // Callers MUST call the returned func, or the subscriber leaks for the life of // the process - one goroutine and one buffered channel per abandoned HTTP // connection, on an endpoint mobile clients reconnect to all day. func (h *Hub) Subscribe() (<-chan string, func()) { // Buffered by one. A doorbell is idempotent - two rings while the listener // is busy mean the same thing as one, because it re-queries from its // cursor either way - so a single slot is enough and a full channel is a // normal state, not backpressure to worry about. ch := make(chan string, 1) h.mu.Lock() id := h.next h.next++ h.subs[id] = ch h.mu.Unlock() return ch, func() { h.mu.Lock() if c, ok := h.subs[id]; ok { delete(h.subs, id) close(c) } h.mu.Unlock() } } // Notify rings every subscriber. Never blocks: a subscriber whose slot is // already full is skipped, because it has a pending wake-up that will make it // re-query anyway. func (h *Hub) Notify(clientID string) { if h == nil || clientID == "" { return } h.mu.Lock() defer h.mu.Unlock() for _, ch := range h.subs { select { case ch <- clientID: default: } } } // Subscribers is the count, for tests and for /healthz. func (h *Hub) Subscribers() int { if h == nil { return 0 } h.mu.Lock() defer h.mu.Unlock() return len(h.subs) }