Files
Behavision/desktop/stream_remote_test.go
Suriyakumarvijayanayagam 48a30d97db Live camera view in the app, and the green light that was lying about it
Two changes, and the second was found by verifying the first.

## Watching a camera from the app, in another building

Snapshots answer "is that camera working". They do not answer "what is
happening in my shop right now", which is what somebody who opens the app away
from the counter is asking. Head office's browser already had that answer -
LiveHub plus cameras.Live, where the shop PC asks outbound whether anybody is
watching and pushes JPEG frames for as long as somebody is - and the app could
not reach it.

cloud.CameraLive opens that feed and the app's own loopback relay re-emits it
as multipart MJPEG. That is the trick: frames arrive base64 over SSE, an <img>
cannot render that, and an <img> renders MJPEG natively - so a tile is an
ordinary <img> pointed at loopback whether the camera is in this room or
another city.

- Reconnecting happens in the relay, not the page. The server caps one push at
  five minutes, so doing it here means the <img> never sees the stream end.
- The headers are flushed before the first frame. Go writes them on the first
  body write, so without that the whole response waits for the shop PC to
  start pushing. Measured against production: 30 seconds and not even a
  Content-Type, which surfaces as the request timing out.
- One camera at a time. Watching makes a shop PC upload, so a grid that went
  live at once would put an estate's worth of cameras on the wire because
  somebody opened a page.
- live.mjpeg is behind the same per-run token as the engine routes, and a
  wrong token is a 404 that never reaches head office at all.
- CameraLive uses its own HTTP client: the shared one's 30s timeout covers the
  whole response and would sever a working view every thirty seconds - the
  trap that made the server set WriteTimeout to zero for its own SSE endpoint.

## A camera read "Connected" for 34 minutes after the shop PC went blind

Which is why the verification above looked like a failure: head office
registered the viewer and no frame ever came.

reportWith returns early when the engine is unreachable - correctly, it has
nothing to say - so the last state it sent stays in the database looking
current. Measured live: cam2 and entrance both reading Connected, in green,
with last_seen_at 34 minutes old, while the heartbeat from the same PC said
cameras_up 0 of 0. Two surfaces reading two stored fields and disagreeing.

false could not be the answer. It means "this camera is not connecting", which
sends an installer to check cabling on a camera that was working perfectly the
last time anybody could ask it. So there are four states and one function:

  connected       reported recently, and working
  not_connecting  reported recently, and the stream will not open
  waiting         no shop PC has ever reported this camera
  stale           reported once, and not lately

- Connected is CLEARED when stale or waiting. A stale true left in place stays
  available to every client reading the field directly, and leaves two fields
  on one object disagreeing - how the shops screen once came out labelled
  Working, in green, above "2 of 3 cameras not connecting".
- Computed in scanCamera, so every camera anybody reads passes through it. A
  state computed per handler is one a handler forgets, and this had already
  reached three screens.
- CameraStaleAfter is 5 minutes: five missed reports, not one. Same reasoning
  as three missed heartbeats - an indicator that cries wolf gets ignored.
- An unparseable last_seen_at is stale. It should be impossible, which is why
  it must not fall through to the state that says everything is fine.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01KGcjxF1cNLcuwc3DAPcnfj
2026-09-30 16:48:43 +05:30

210 lines
6.8 KiB
Go

package main
import (
"bytes"
"context"
"encoding/base64"
"fmt"
"io"
"net/http"
"net/http/httptest"
"strings"
"sync/atomic"
"testing"
"time"
)
// jpg is a byte sequence that is not valid JPEG and does not need to be: what
// is under test is that the bytes arrive intact and framed, not that a decoder
// likes them.
var jpg = []byte{0xFF, 0xD8, 'h', 'e', 'l', 'l', 'o', 0xFF, 0xD9}
// sseServer answers head office's live endpoint with `pushes` frames and then
// ends the response, which is what the server's five-minute cap does.
func sseServer(t *testing.T, frames int, hits *int32) *httptest.Server {
t.Helper()
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
atomic.AddInt32(hits, 1)
w.Header().Set("Content-Type", "text/event-stream")
fl, _ := w.(http.Flusher)
// Registered, nothing being pushed yet. Nothing may be drawn for it.
fmt.Fprint(w, "event: waiting\ndata: \n\n")
if fl != nil {
fl.Flush()
}
for i := 0; i < frames; i++ {
fmt.Fprintf(w, "event: frame\ndata: %s\n\n",
base64.StdEncoding.EncodeToString(jpg))
if fl != nil {
fl.Flush()
}
}
}))
}
func openerFor(srv *httptest.Server) func(context.Context, string) (*http.Response, error) {
return func(ctx context.Context, cam string) (*http.Response, error) {
req, _ := http.NewRequestWithContext(ctx, http.MethodGet, srv.URL+"/live/"+cam, nil)
return http.DefaultClient.Do(req)
}
}
// The whole point: base64 frames over SSE are not something an <img> can show,
// and a multipart MJPEG stream is. Without this the app could only ever show a
// still, on exactly the computers that cannot reach the camera any other way.
func TestRemoteFramesReachTheWebviewAsMJPEG(t *testing.T) {
var hits int32
srv := sseServer(t, 3, &hits)
defer srv.Close()
p := newStreamProxy()
if err := p.watchRemote(openerFor(srv)); err != nil {
t.Fatalf("watchRemote: %v", err)
}
defer p.stop()
u := p.urlFor("cam2", "live.mjpeg")
if u == "" {
t.Fatal("no relay url; the proxy did not bind")
}
// The relay reconnects for as long as the viewer is there, so the read is
// bounded by us rather than by the stream ending - exactly as an <img>
// would behave.
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodGet, u, nil)
resp, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatalf("GET relay: %v", err)
}
defer resp.Body.Close()
if ct := resp.Header.Get("Content-Type"); !strings.HasPrefix(ct, "multipart/x-mixed-replace") {
t.Fatalf("Content-Type = %q, an <img> will not treat that as a stream", ct)
}
// Read the first three frames' worth and stop; the relay would otherwise
// go on reconnecting forever, which is the behaviour being relied on.
want := append([]byte(fmt.Sprintf("--%s\r\nContent-Type: image/jpeg\r\nContent-Length: %d\r\n\r\n",
mjpegBoundary, len(jpg))), jpg...)
got := make([]byte, 0, 4096)
buf := make([]byte, 512)
for len(got) < 3*len(want) {
n, rerr := resp.Body.Read(buf)
got = append(got, buf[:n]...)
if rerr != nil {
break
}
}
if n := bytes.Count(got, []byte("--"+mjpegBoundary)); n < 3 {
t.Fatalf("got %d frames in %d bytes, want at least 3", n, len(got))
}
if !bytes.Contains(got, want) {
t.Errorf("a frame was not framed as expected:\n%q", got[:min(len(got), 300)])
}
// `waiting` is a real state - head office has us registered and the shop
// computer has not started pushing - and there is nothing to draw for it.
// Emitting an empty part would blank a tile that already had a picture.
if bytes.Contains(got, []byte("Content-Length: 0")) {
t.Error("an empty frame was written for a waiting event")
}
}
// The server caps one push at five minutes so a tab left open for a week
// cannot leave a shop uploading for a week. Reconnecting is therefore a normal
// event, and doing it here rather than in the page is what lets the <img>
// survive the cap - it never sees the stream end.
func TestTheRelayReconnectsWhenHeadOfficeEndsAPush(t *testing.T) {
var hits int32
srv := sseServer(t, 1, &hits)
defer srv.Close()
p := newStreamProxy()
if err := p.watchRemote(openerFor(srv)); err != nil {
t.Fatalf("watchRemote: %v", err)
}
defer p.stop()
ctx, cancel := context.WithTimeout(context.Background(), 4*time.Second)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodGet, p.urlFor("cam2", "live.mjpeg"), nil)
resp, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatalf("GET relay: %v", err)
}
defer resp.Body.Close()
// Two frames means two pushes, because each push carries exactly one.
seen, buf := 0, make([]byte, 256)
acc := make([]byte, 0, 2048)
for seen < 2 {
n, rerr := resp.Body.Read(buf)
acc = append(acc, buf[:n]...)
seen = bytes.Count(acc, []byte("--"+mjpegBoundary))
if rerr != nil {
break
}
}
if seen < 2 {
t.Fatalf("got %d frames across reconnects, want 2", seen)
}
if got := atomic.LoadInt32(&hits); got < 2 {
t.Errorf("head office was asked %d times, want at least 2", got)
}
}
// Signed out, the relay must not pretend. There is no fallback URL to offer
// either: the head-office endpoint needs this session's bearer, which an <img>
// cannot send - so a tile that silently failed would be the only alternative.
func TestTheRelayRefusesWhenNobodyIsSignedIn(t *testing.T) {
p := newStreamProxy()
if err := p.watchRemote(nil); err != nil {
t.Fatalf("watchRemote: %v", err)
}
defer p.stop()
resp, err := http.Get(p.urlFor("cam2", "live.mjpeg"))
if err != nil {
t.Fatalf("GET relay: %v", err)
}
defer resp.Body.Close()
io.Copy(io.Discard, resp.Body)
if resp.StatusCode != http.StatusBadGateway {
t.Errorf("status = %d, want 502", resp.StatusCode)
}
}
// The relay is credentialed - it is a path to a live view of a shop floor -
// and the token is the only thing standing between another local process and
// it. live.mjpeg must be behind exactly the same door as the engine routes.
func TestTheRemoteRouteIsBehindTheSameToken(t *testing.T) {
var hits int32
srv := sseServer(t, 1, &hits)
defer srv.Close()
p := newStreamProxy()
if err := p.watchRemote(openerFor(srv)); err != nil {
t.Fatalf("watchRemote: %v", err)
}
defer p.stop()
// The right shape, the wrong value.
parts := strings.Split(p.urlFor("cam2", "live.mjpeg"), "/")
parts[4] = strings.Repeat("0", len(parts[4]))
bad := strings.Join(parts, "/")
resp, err := http.Get(bad)
if err != nil {
t.Fatalf("GET relay: %v", err)
}
defer resp.Body.Close()
io.Copy(io.Discard, resp.Body)
if resp.StatusCode != http.StatusNotFound {
t.Errorf("status = %d, want 404 - and 404 rather than 403, because there is nothing here to tell an unwelcome caller they found the right door", resp.StatusCode)
}
if atomic.LoadInt32(&hits) != 0 {
t.Error("a request with the wrong token still made the shop computer upload")
}
}