// Command behavision-server consumes store events from the broker and serves // the HTTP API. // // One process, two jobs, because they share the database pool and the volume // does not justify splitting them. If ingest ever needs to scale separately it // can: nothing in `ingest` knows about HTTP. package main import ( "context" "encoding/json" "errors" "fmt" "log" "net/http" "os" "os/signal" "strings" "sync/atomic" "syscall" "time" paho "github.com/eclipse/paho.mqtt.golang" "github.com/loyaly/behavision-server/internal/api" "github.com/loyaly/behavision-server/internal/assistant" "github.com/loyaly/behavision-server/internal/blob" "github.com/loyaly/behavision-server/internal/broker" "github.com/loyaly/behavision-server/internal/ingest" "github.com/loyaly/behavision-server/internal/migrate" "github.com/loyaly/behavision-server/internal/secret" "github.com/loyaly/behavision-server/internal/store" "github.com/loyaly/behavision-server/internal/web" "github.com/loyaly/behavision-server/migrations" ) var version = "dev" func main() { // One binary, two jobs. Provisioning needs the same database URL and the // same secret key as the server, and a second image to keep in sync with // this one is a second thing to forget to deploy. if len(os.Args) > 1 && os.Args[1] == "provision" { if err := runProvision(os.Args[2:]); err != nil { fmt.Fprintln(os.Stderr, err) os.Exit(1) } return } if len(os.Args) > 1 && os.Args[1] == "broker-init" { if err := runBrokerInit(os.Args[2:]); err != nil { fmt.Fprintln(os.Stderr, err) os.Exit(1) } return } if len(os.Args) > 1 && os.Args[1] == "migrate" { if err := runMigrate(os.Args[2:]); err != nil { fmt.Fprintln(os.Stderr, err) os.Exit(1) } return } if err := run(); err != nil { log.Fatalf("behavision-server: %v", err) } } func env(key, def string) string { if v := os.Getenv(key); v != "" { return v } return def } func run() error { logger := log.New(os.Stdout, "", log.LstdFlags|log.LUTC) dsn := os.Getenv("DATABASE_URL") if dsn == "" { return errors.New("DATABASE_URL is required") } brokerURL := env("MQTT_URL", "tcp://behavision-mqtt:1883") brokerUser := os.Getenv("MQTT_USERNAME") brokerPass := os.Getenv("MQTT_PASSWORD") addr := env("LISTEN_ADDR", ":8080") ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() st, err := store.Open(ctx, dsn) if err != nil { return err } defer st.Close() st.UseLogger(logger) logger.Print("database connected") // Applied here, not by hand, because an upgrade of this product is "copy // the new binary and restart it". A migration an operator has to remember // to run is a migration that does not get run, and its failure is not a // startup error - it is a query referencing a column that is not there, // surfacing later on whichever endpoint touches it first. // // Refusing to start on failure is deliberate: a server running against a // schema it does not match writes wrong data, and wrong data outlives the // outage that stopping causes. if os.Getenv("BEHAVISION_SKIP_MIGRATE") == "1" { logger.Print("WARN BEHAVISION_SKIP_MIGRATE=1 - schema not checked") } else { applied, err := migrate.Apply(ctx, st.Pool(), migrations.FS) if err != nil { return fmt.Errorf("schema: %w", err) } if len(applied) == 0 { logger.Print("schema up to date") } else { logger.Printf("schema: applied %s", strings.Join(applied, ", ")) } } // Without the key the server still ingests events and serves reports; only // enrolment fails, and it fails with a message naming the missing variable. // Refusing to start would take a working estate down over a feature that // runs once per shop PC. if box, err := secret.FromEnv("BEHAVISION_SECRET_KEY"); err != nil { logger.Printf("WARN %v - agent enrolment will be refused", err) } else { st.UseSecrets(box) } // One hub, shared by the MQTT consumer and the API in this single process. // The consumer rings it; live arrival streams answer by re-querying. If // this ever runs as more than one instance, an instance will not hear the // others' ingest and its streams fall back to their slow tick - latency, // not silence, which is why the fallback exists at all. hub := api.NewHub() consumer := &ingest.Consumer{Store: st, Log: logger, Notify: hub.Notify} var brokerUp atomic.Bool opts := paho.NewClientOptions(). AddBroker(brokerURL). SetClientID(fmt.Sprintf("behavision-server-%d", time.Now().UnixNano())). SetUsername(brokerUser). SetPassword(brokerPass). SetAutoReconnect(true). SetConnectRetry(true). SetConnectRetryInterval(5 * time.Second). SetKeepAlive(30 * time.Second). // Messages are handled one at a time in order. Concurrent handlers // would let two events for the same new visitor race and create two // visitor rows for one person. SetOrderMatters(true) opts.OnConnectionLost = func(_ paho.Client, err error) { brokerUp.Store(false) logger.Printf("broker connection lost: %v", err) } opts.OnConnect = func(c paho.Client) { brokerUp.Store(true) // Subscribing inside OnConnect, not once after Connect: a reconnect // with a clean session drops server-side subscriptions, and without // this the process would sit there connected and deaf. tok := c.Subscribe("bv/#", 1, func(_ paho.Client, m paho.Message) { hctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) defer cancel() if err := consumer.Handle(hctx, m.Topic(), m.Payload()); err != nil { // Returning without acking is not possible with this client's // auto-ack, so a transient failure is logged loudly rather than // silently losing the event. QoS 1 + the agent's own spool mean // the event still exists at the far end until we confirm it. logger.Printf("ERROR handling %s: %v", m.Topic(), err) } }) if tok.Wait() && tok.Error() != nil { logger.Printf("subscribe failed: %v", tok.Error()) return } logger.Printf("subscribed to bv/# on %s", brokerURL) } client := paho.NewClient(opts) // Connected in the BACKGROUND, and that is the whole point of the comment // below: the API must come up even when the broker is down. Waiting here // did the opposite. With SetConnectRetry the token does not complete until // the broker answers, so an unreachable broker held the HTTP listener down // for the full 20 seconds on every start - measured, on a machine with no // broker at all. And it was silent: WaitTimeout returns false on a timeout, // which short-circuits the && , so the one log line never printed either. go func() { tok := client.Connect() if !tok.WaitTimeout(20 * time.Second) { logger.Printf("broker %s not reachable yet - retrying in the background", brokerURL) return } if err := tok.Error(); err != nil { logger.Printf("initial broker connect failed (will retry): %v", err) } }() apiSrv := &api.Server{ Store: st, Log: logger, // Opening a shop registers its broker login at the same moment, over the // same broker credential the ingest side already holds. Broker: broker.New(brokerURL, brokerUser, brokerPass, logger), Blob: objectStore(ctx, logger), Hub: hub, Bootstrap: api.BootstrapConfig{ // What an enrolling PC is told to connect to. From the server's own // environment, never from the request: an agent asking where to // connect must not get to influence the answer. MQTTURL: env("AGENT_MQTT_URL", "tls://mcp.loyaly.ai:8883"), CACert: caCert(logger), Models: modelManifest(logger), }, } // The assistant. Its tools are the same business questions the screens ask, // and every one runs as the signed-in user. Without ANTHROPIC_API_KEY it // reports itself off and the UI hides the panel - a supported state, not a // startup failure, because everything else on this server still works. assistantTools := &assistant.Registry{ Store: st, SiteChecker: api.BuildSiteSteps, Now: func() time.Time { return time.Now().UTC() }, } assistantClient := &assistant.Client{Tools: assistantTools, Log: logger} if assistantClient.Configured() { apiSrv.Assistant = assistant.ForAPI(assistantClient) logger.Print("assistant enabled") } else { logger.Print("ANTHROPIC_API_KEY not set - the assistant is off " + "(everything else is unaffected)") } mux := apiSrv.Routes() mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { hctx, cancel := context.WithTimeout(r.Context(), 3*time.Second) defer cancel() dbErr := st.Ping(hctx) body := map[string]any{ "version": version, "database": dbErr == nil, "broker": brokerUp.Load(), "accepted": consumer.Accepted, "duplicate": consumer.Duplicate, "dropped": consumer.Dropped, } // Degraded, not down: the database is the hard dependency. Reporting // unhealthy while the broker reconnects would make an orchestrator // restart a process that is working. if dbErr != nil { body["error"] = dbErr.Error() w.WriteHeader(http.StatusServiceUnavailable) } w.Header().Set("Content-Type", "application/json") json.NewEncoder(w).Encode(body) }) // Anything under /api that matched no route stays JSON. Falling through to // the web app would turn a typo'd endpoint into an HTML page arriving where // a client expects JSON, and the parse error it causes surfaces three // layers from the cause. mux.HandleFunc("/api/", func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") w.WriteHeader(http.StatusNotFound) fmt.Fprint(w, `{"error":"not_found","message":"No such endpoint."}`) }) // Everything else is the head-office web app, embedded in this binary. site, err := web.Handler() if err != nil { return fmt.Errorf("web assets: %w", err) } mux.Handle("/", site) srv := &http.Server{ Addr: addr, Handler: securityHeaders(mux), // Bounded so a slow or hostile client cannot hold a connection open // indefinitely on a 2 vCPU box. ReadHeaderTimeout: 10 * time.Second, ReadTimeout: 30 * time.Second, // WriteTimeout is deliberately ZERO, and the arrivals stream is why. // // A write deadline applies to the WHOLE response, not to each write, so // any non-zero value silently severs every event stream that outlives // it - a shop screen that dies after 60 seconds and reconnects forever, // which looks like a network fault and is not one. ReadTimeout and // IdleTimeout still bound a slow or hostile client; what is given up is // a cap on how long a client may take to READ a response it asked for. WriteTimeout: 0, IdleTimeout: 120 * time.Second, } go func() { logger.Printf("http listening on %s", addr) if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { logger.Printf("http server stopped: %v", err) stop() } }() <-ctx.Done() logger.Print("shutting down") shutCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() srv.Shutdown(shutCtx) //nolint:errcheck client.Disconnect(1000) return nil } func securityHeaders(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { // Traefik terminates TLS, so HSTS is set here where the app knows it is // only ever served over https. w.Header().Set("Strict-Transport-Security", "max-age=31536000") w.Header().Set("X-Content-Type-Options", "nosniff") w.Header().Set("Referrer-Policy", "no-referrer") w.Header().Set("X-Frame-Options", "DENY") if strings.HasPrefix(r.URL.Path, "/api/") { w.Header().Set("Cache-Control", "no-store") } next.ServeHTTP(w, r) }) } // caCert reads the CA an agent needs to trust the broker's certificate. // // Handed out at enrolment rather than shipped in the installer: the CA can be // rotated without re-signing and re-distributing every store's software. func caCert(logger *log.Logger) string { path := os.Getenv("AGENT_CA_FILE") if path == "" { return "" } b, err := os.ReadFile(path) if err != nil { logger.Printf("WARN cannot read AGENT_CA_FILE %s: %v - "+ "enrolling agents will get no CA and cannot verify the broker", path, err) return "" } return string(b) } // modelManifest pins which encoder new installs download. // // It matters more than it looks: embeddings are model-tagged and vectors from // two different encoders cannot be compared at all, so a fleet that drifts onto // two models is a fleet whose sites cannot recognise each other's customers. func modelManifest(logger *log.Logger) []api.ModelRef { path := os.Getenv("AGENT_MODELS_FILE") if path == "" { return nil } b, err := os.ReadFile(path) if err != nil { logger.Printf("WARN cannot read AGENT_MODELS_FILE %s: %v", path, err) return nil } var out []api.ModelRef if err := json.Unmarshal(b, &out); err != nil { logger.Printf("WARN AGENT_MODELS_FILE is not valid JSON: %v", err) return nil } return out } // objectStore configures face-image storage, or returns nil for a deployment // that stores none. // // Nil is a supported, and the default, configuration: this product shipped // deliberately storing no images at all, and switching them on changes what the // database is under GDPR and India's DPDP. It should take setting five // variables, not forgetting to unset them. // // The self-test is not optional. The bucket in use is world-readable at the // bucket level, so "the upload worked" and "the face image is private" are // different questions, and the second one has to be answered at boot rather // than discovered later from someone else's search results. func objectStore(ctx context.Context, logger *log.Logger) api.BlobStore { cfg := blob.Config{ Region: os.Getenv("DO_SPACES_REGION"), Endpoint: os.Getenv("DO_SPACES_ENDPOINT"), Bucket: os.Getenv("DO_SPACES_BUCKET"), AccessKey: os.Getenv("DO_SPACES_ACCESS_KEY"), SecretKey: os.Getenv("DO_SPACES_SECRET_KEY"), Prefix: env("DO_SPACES_PREFIX", "behavision/v2"), } if cfg.AccessKey == "" && cfg.Bucket == "" { logger.Print("object storage not configured - visits will carry no images") return nil } store, err := blob.New(cfg) if err != nil { logger.Printf("WARN object storage disabled: %v", err) return nil } cctx, cancel := context.WithTimeout(ctx, 20*time.Second) defer cancel() if err := store.Check(cctx); err != nil { // Refuse to serve images rather than serve them unsafely. A shop losing // photos is a visible, fixable problem; a shop publishing its // customers' faces is neither. logger.Printf("ERROR object storage self-test failed, images DISABLED: %v", err) return nil } logger.Printf("object storage ready: %s/%s (objects are private)", cfg.Bucket, cfg.Prefix) return store }