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
400 lines
14 KiB
Go
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
|
|
}
|