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
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
// 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
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
// 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
// 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")
}
}