// End-to-end probe: publish visits through the real agent path. // // SIM controls the similarity of the second visit to the first, which is what // exercises the reinforcement branch: identical vectors teach the gallery // nothing and must be refused, a genuinely different view of the same person // must be kept. package main import ( "context" "encoding/json" "fmt" "log" "math" "math/rand" "os" "strconv" "time" "github.com/loyaly/behavision-agent/pkg/mqtt" "github.com/loyaly/behavision-agent/pkg/spool" ) func unit(seed int64) []float32 { rng := rand.New(rand.NewSource(seed)) v := make([]float32, 512) var n float64 for i := range v { v[i] = float32(rng.NormFloat64()) n += float64(v[i]) * float64(v[i]) } n = math.Sqrt(n) for i := range v { v[i] /= float32(n) } return v } // atSimilarity builds a unit vector exactly `target` from base. func atSimilarity(base []float32, target float64, seed int64) []float32 { other := unit(seed) var dot float64 for i := range base { dot += float64(base[i]) * float64(other[i]) } var n float64 for i := range other { other[i] -= float32(dot) * base[i] n += float64(other[i]) * float64(other[i]) } n = math.Sqrt(n) out := make([]float32, len(base)) k := math.Sqrt(1 - target*target) for i := range base { out[i] = float32(target)*base[i] + float32(k)*(other[i]/float32(n)) } return out } func main() { broker, user := os.Getenv("BROKER"), os.Getenv("MQTT_USER") logger := log.New(os.Stdout, " ", 0) q, err := spool.Open(os.Getenv("SPOOL_DIR"), 100) if err != nil { log.Fatal(err) } sim, _ := strconv.ParseFloat(os.Getenv("SIM"), 64) emb := unit(42) if sim > 0 { emb = atSimilarity(unit(42), sim, 7) } quality, _ := strconv.ParseFloat(os.Getenv("QUALITY"), 64) if quality == 0 { quality = 0.74 } visit := map[string]any{ "event_id": os.Getenv("EVENT_ID"), "occurred_at": time.Now().UTC().Format(time.RFC3339Nano), "camera_id": "entrance", "is_new": sim == 0, "quality": quality, "model": "w600k_r50.onnx", "embedding": emb, "attributes": map[string]any{"gender": "Male", "age": 34}, } if err := q.Append(fmt.Sprintf("bv/%s/visit", user), visit); err != nil { log.Fatal(err) } client, err := mqtt.NewClient(mqtt.ClientOptions{ BrokerURL: broker, ClientID: "e2e-probe-" + os.Getenv("EVENT_ID"), Username: user, Password: os.Getenv("MQTT_PASS"), CAFile: os.Getenv("CA_FILE"), Log: logger, }) if err != nil { log.Fatalf("connect: %v", err) } defer client.Close() ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) defer cancel() go (&mqtt.Pump{Queue: q, Publisher: client, Log: logger}).Run(ctx) deadline := time.Now().Add(15 * time.Second) for time.Now().Before(deadline) { if q.Len() == 0 { b, _ := json.Marshal(map[string]any{"sim": sim, "quality": quality}) logger.Printf("published %s", b) return } time.Sleep(200 * time.Millisecond) } log.Fatalf("spool did not drain: %d left", q.Len()) }