Files
Behavision/server/cmd/behavision-server/main.go
Suriyakumarvijayanayagam 5453c26e4c The schema applies itself, and the setup script stops hiding failures
Migrations were run by hand and nothing recorded which had run, so
re-running the setup script against an existing database failed on the
first CREATE TABLE, and shipping a new migration gave an operator no way
to know whether an estate had it. A missed migration is not a startup
error - it is a query referencing a column that is not there, surfacing
later on whichever endpoint touches it first.

server/internal/migrate applies pending migrations at boot and refuses to
start against a schema it does not match. One transaction per file
holding both the DDL and the row that records it; an advisory lock so two
servers starting at once cannot both apply 008; checksums so an edited
migration is refused by name rather than silently skipped; numeric
ordering so 010 does not run before 009. `migrate -baseline N` adopts a
database built before any of this existed, because "the clients table
exists" does not say whether 007's index does.

Verified on the live database: adopted 001-007, applied 008.

008 adds two indexes on `purchases`, found by asking the database which
foreign keys had nothing behind them and then checking what queries the
table. The conversion report filters client_id + occurred_at, which is
exactly the estate-wide case with no site to narrow it.

run-local.sh had two bugs, both found by running it rather than reading
it: it reused a broker container whose bind mount pointed at a directory
that no longer existed, and it discarded stderr on the mosquitto_passwd
call, so under `set -e` it exited at step 5 with no output at all.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HViLj9gYNRtSr7YVZmW5sn
2026-09-04 12:06:52 +05:30

400 lines
14 KiB
Go

// 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/ingest"
"github.com/loyaly/behavision-server/internal/migrate"
"github.com/loyaly/behavision-server/internal/web"
"github.com/loyaly/behavision-server/internal/secret"
"github.com/loyaly/behavision-server/internal/store"
"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] == "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,
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
}