package engine import ( "bytes" "context" "net/http" "net/http/httptest" "os/exec" "sync" "testing" "time" ) // sh supervises /bin/sh instead of a 200 MB frozen engine. The Command hook // exists for exactly this. func sh(script string) func(context.Context) *exec.Cmd { return func(ctx context.Context) *exec.Cmd { return exec.CommandContext(ctx, "/bin/sh", "-c", script) } } func waitFor(t *testing.T, s *Supervisor, want State, within time.Duration) { t.Helper() deadline := time.Now().Add(within) for time.Now().Before(deadline) { if got, _ := s.State(); got == want { return } time.Sleep(5 * time.Millisecond) } got, err := s.State() t.Fatalf("state %q (err %v), want %q within %s", got, err, want, within) } func TestItRunsAndReportsRunning(t *testing.T) { s := New(Options{Command: sh("sleep 5")}) s.Start() defer s.Stop() waitFor(t, s, Running, 2*time.Second) } func TestStopDoesNotTriggerARestart(t *testing.T) { // The classic supervisor bug: the user presses Stop, the child exits, the // loop reads that as a crash and starts it again. s := New(Options{Command: sh("sleep 30")}) s.Start() waitFor(t, s, Running, 2*time.Second) s.Stop() if got, _ := s.State(); got != Stopped { t.Fatalf("state after Stop is %q", got) } if n := s.Restarts(); n != 0 { t.Fatalf("Stop counted as %d crash-restarts", n) } time.Sleep(200 * time.Millisecond) if got, _ := s.State(); got != Stopped { t.Fatalf("it restarted itself after Stop: %q", got) } } func TestStopIsSynchronous(t *testing.T) { // Stop must not return while the child still holds the SQLite WAL, or the // next Start races the previous process. s := New(Options{Command: sh("sleep 30")}) s.Start() waitFor(t, s, Running, 2*time.Second) done := make(chan struct{}) go func() { s.Stop(); close(done) }() select { case <-done: case <-time.After(3 * time.Second): t.Fatal("Stop did not return") } } func TestACrashIsRestarted(t *testing.T) { s := New(Options{Command: sh("exit 1")}) s.Start() defer s.Stop() deadline := time.Now().Add(3 * time.Second) for time.Now().Before(deadline) { if s.Restarts() >= 2 { return } time.Sleep(10 * time.Millisecond) } t.Fatalf("only %d restarts - is it backing off correctly?", s.Restarts()) } func TestItDoesNotSpinOnAProcessThatCannotStart(t *testing.T) { // A tight restart loop on a broken install pins a core and fills the disk // with log lines. Backoff must space the attempts out. s := New(Options{Command: sh("exit 1")}) s.Start() defer s.Stop() time.Sleep(1500 * time.Millisecond) // 1s + 2s backoff means at most ~2 attempts in 1.5s; a spin would be // thousands. if n := s.Restarts(); n > 4 { t.Fatalf("%d restarts in 1.5s - not backing off", n) } } func TestItGivesUpAfterMaxRestarts(t *testing.T) { s := New(Options{Command: sh("exit 1"), MaxRestarts: 2}) s.Start() defer s.Stop() waitFor(t, s, Failed, 5*time.Second) if _, err := s.State(); err == nil { t.Fatal("Failed state carries no reason") } } func TestEngineOutputIsCaptured(t *testing.T) { // A crashed engine with no captured output means a site visit to diagnose. var mu sync.Mutex buf := &lockedBuf{mu: &mu} s := New(Options{Command: sh("echo model-load-failed; exit 1"), LogWriter: buf, MaxRestarts: 1}) s.Start() defer s.Stop() waitFor(t, s, Failed, 5*time.Second) if got := buf.String(); !bytes.Contains([]byte(got), []byte("model-load-failed")) { t.Fatalf("engine output not captured, got %q", got) } } func TestStartTwiceDoesNotRunTwoEngines(t *testing.T) { // Two engines on one SQLite WAL and one camera is the failure this whole // package exists to prevent. s := New(Options{Command: sh("sleep 5")}) s.Start() s.Start() defer s.Stop() waitFor(t, s, Running, 2*time.Second) if n := s.Restarts(); n != 0 { t.Fatalf("second Start disturbed the first: %d restarts", n) } } func TestHealthReportsTheModelThatActuallyLoaded(t *testing.T) { // A running process is not a working engine: on a memory-starved box the // big model loses the fallback chain and the process stays up regardless. srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { user, pass, ok := r.BasicAuth() if !ok || user != "u" || pass != "p" { w.WriteHeader(http.StatusUnauthorized) return } w.Write([]byte(`{"status":"ok","recognition_model":"w600k_mbf.onnx", "cameras":{"entrance":true}}`)) })) defer srv.Close() s := New(Options{Command: sh("sleep 1"), HealthURL: srv.URL, User: "u", Password: "p"}) h, err := s.Health(context.Background()) if err != nil { t.Fatal(err) } if h.RecognitionModel != "w600k_mbf.onnx" || !h.Cameras["entrance"] { t.Fatalf("bad health: %+v", h) } } func TestHealthFailsClosedOnBadCredentials(t *testing.T) { srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusUnauthorized) })) defer srv.Close() s := New(Options{Command: sh("true"), HealthURL: srv.URL, User: "u", Password: "wrong"}) if _, err := s.Health(context.Background()); err == nil { t.Fatal("401 reported as healthy") } } type lockedBuf struct { mu *sync.Mutex buf bytes.Buffer } func (l *lockedBuf) Write(p []byte) (int, error) { l.mu.Lock() defer l.mu.Unlock() return l.buf.Write(p) } func (l *lockedBuf) String() string { l.mu.Lock() defer l.mu.Unlock() return l.buf.String() } // The engine has to be TOLD where to post detections, and the only place that // can happen is when the child is launched: the bridge picks a random loopback // port after the supervisor is built, and a restarted engine has to be told // again. This pins that the Command hook is consulted per launch rather than // captured once - the wiring that was missing while the bridge's own doc // comment claimed it existed. func TestTheChildIsBuiltFreshOnEveryLaunch(t *testing.T) { var mu sync.Mutex url := "http://127.0.0.1:1111/e" var seen []string s := New(Options{Command: func(ctx context.Context) *exec.Cmd { mu.Lock() seen = append(seen, url) mu.Unlock() return exec.CommandContext(ctx, "/bin/sh", "-c", "exit 1") }}) s.Start() waitFor(t, s, Backoff, 2*time.Second) // The port changes, exactly as it does when the bridge restarts. mu.Lock() url = "http://127.0.0.1:2222/e" mu.Unlock() // One backoff (1s) plus room for the relaunch. deadline := time.Now().Add(4 * time.Second) for time.Now().Before(deadline) { mu.Lock() n := len(seen) mu.Unlock() if n >= 2 { break } time.Sleep(20 * time.Millisecond) } s.Stop() mu.Lock() defer mu.Unlock() if len(seen) < 2 { t.Fatalf("the command hook ran %d times, so a restart could not be told a new URL", len(seen)) } if seen[len(seen)-1] != "http://127.0.0.1:2222/e" { t.Fatalf("the last launch used %q - the hook captured a stale value", seen[len(seen)-1]) } }