// Package ingest turns broker messages into database rows. // // The whole tenancy decision happens here, in one place, so it can be read and // tested as a unit rather than being spread across handlers. package ingest import ( "context" "encoding/json" "errors" "fmt" "log" "github.com/loyaly/behavision-server/internal/contract" ) // Store is the persistence the consumer needs. An interface so the routing and // tenancy rules below can be tested without a database. type Store interface { // ResolveSite maps an authenticated MQTT username to a tenant. It must // NEVER create anything: see the comment in Handle. ResolveSite(ctx context.Context, mqttUsername string) (Site, error) RecordVisit(ctx context.Context, site Site, v *contract.Visit) (inserted bool, err error) RecordHeartbeat(ctx context.Context, site Site, h *contract.Heartbeat) error } // Site is a resolved tenant. type Site struct { ClientID string SiteID string AgentID string Slug string } // ErrUnknownSite means no provisioned site matches the username. var ErrUnknownSite = errors.New("unknown site") // Consumer applies one message at a time. type Consumer struct { Store Store Log *log.Logger // Metrics, read by /healthz. Counting drops matters as much as counting // successes: a consumer silently discarding a tenth of its traffic looks // identical to a quiet week. Accepted uint64 Duplicate uint64 Dropped uint64 // Notify, when set, is rung once per visit that was genuinely new, so live // listeners re-query instead of waiting out their fallback tick. It carries // a client id and nothing else - the listeners run the same query a polling // client would, so there is one definition of what an arrival looks like. // // It MUST NOT block. This runs on the consumer's own goroutine, and one // slow subscriber holding it up would stall ingest for every tenant on the // server, turning a cosmetic delay on one shop screen into lost throughput // across the estate. Notify func(clientID string) } // Handle processes one broker message. // // Returning nil means "done with this message, do not redeliver". A permanent // failure returns nil too, on purpose: retrying a malformed payload forever // would stop every good message behind it, which is exactly the failure the // agent's spool quarantine exists to avoid. Only a transient error asks for // redelivery. func (c *Consumer) Handle(ctx context.Context, topic string, payload []byte) error { t, err := contract.ParseTopic(topic) if err != nil { c.drop("topic %q: %v", topic, err) return nil } // The tenant comes from the topic prefix, which the broker enforces with // `pattern write bv/%u/...` - a site physically cannot publish under // another site's prefix. So this lookup is a resolution, not a trust // decision. // // It must never auto-create. A site typo'd into existence would silently // become a tenant with its own visitors and its own footfall report, and // nobody would notice until the numbers did not add up. Provisioning is a // deliberate act. site, err := c.Store.ResolveSite(ctx, t.Username) if err != nil { if errors.Is(err, ErrUnknownSite) { c.drop("no provisioned site for %q (broker credential exists but "+ "the site was never registered)", t.Username) return nil } return fmt.Errorf("resolve site %q: %w", t.Username, err) } switch t.Kind { case "visit": return c.handleVisit(ctx, site, payload) case "heartbeat": return c.handleHeartbeat(ctx, site, payload) case "status": // Same shape as a heartbeat, sent on change rather than on a timer. return c.handleHeartbeat(ctx, site, payload) default: c.drop("unknown message kind %q from %s", t.Kind, t.Username) return nil } } func (c *Consumer) handleVisit(ctx context.Context, site Site, payload []byte) error { var v contract.Visit if err := json.Unmarshal(payload, &v); err != nil { c.drop("visit from %s is not valid json: %v", site.Slug, err) return nil } if err := v.Validate(); err != nil { if errors.Is(err, contract.ErrPermanent) { c.drop("visit %q from %s rejected: %v", v.EventID, site.Slug, err) return nil } return err } inserted, err := c.Store.RecordVisit(ctx, site, &v) if err != nil { // A database failure IS transient - the broker should redeliver rather // than the event being lost. This is the one path that returns an error. return fmt.Errorf("record visit %q: %w", v.EventID, err) } if inserted { c.Accepted++ // Only on a genuine insert. At-least-once delivery makes redelivery // normal, and ringing for a duplicate would wake every live stream on // the estate to re-query for a row they already have. if c.Notify != nil { c.Notify(site.ClientID) } } else { // Not an error and not a warning: at-least-once delivery makes this // normal after any reconnect. It is counted so that a sudden rise - // which would mean the agent is not acking - is visible. c.Duplicate++ } return nil } func (c *Consumer) handleHeartbeat(ctx context.Context, site Site, payload []byte) error { var h contract.Heartbeat if err := json.Unmarshal(payload, &h); err != nil { c.drop("heartbeat from %s is not valid json: %v", site.Slug, err) return nil } if err := c.Store.RecordHeartbeat(ctx, site, &h); err != nil { return fmt.Errorf("record heartbeat for %s: %w", site.Slug, err) } if h.Dropped > 0 { // The site's own queue overflowed, so it genuinely lost footfall. // Surfaced loudly - this is data the customer paid for and will never // get back, and it must not be discoverable only by reading a graph. c.logf("WARNING: site %s reports %d events dropped from its local "+ "queue (it was offline long enough to overflow)", site.Slug, h.Dropped) } c.Accepted++ return nil } func (c *Consumer) drop(format string, args ...any) { c.Dropped++ c.logf("dropped: "+format, args...) } func (c *Consumer) logf(format string, args ...any) { if c.Log != nil { c.Log.Printf(format, args...) } }