Asked of the row the feed actually returns. site_id had a reference all along and the feed was not sending it. A client could read the shop's NAME off an arrival and still had no way to ask for that shop except by uuid - the exact gap the reference scheme exists to close. site_slug now travels with it. visit_id stays a uuid and needs no reference: no route takes it, it is a key a client de-duplicates on because delivery is at-least-once, and nobody says a visit id out loud. The uuid in a face URL must STAY random. visit_faces.id is gen_random_uuid() and a derived or sequential one would let somebody walk a shop's customers by date - the same reason bucket keys are random rather than derived from the event id. A readable identifier is right for a customer and wrong for the thing that points at their photograph. And seq is now json:"-". visits.seq is a plain bigserial, so it counts every visit on the PLATFORM, and shipping it put the total footfall of every customer we have on every row of every tenant's feed - the same German-tank estimate that decided visitors.number had to be per client. It was a convenience for "have I fallen behind", nothing ever read it, and the cursor answers that without disclosing a number. The SSE event id was never the raw value; it has always been the opaque cursor. The one test that broke was reading seq back off the wire to assert the cursor pointed at the last row of a burst. It asserts against the seeded position now: the property is unchanged, and the test can no longer see what a client cannot. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HViLj9gYNRtSr7YVZmW5sn
700 lines
22 KiB
Go
700 lines
22 KiB
Go
package api
|
|
|
|
import (
|
|
"context"
|
|
"encoding/base64"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"strings"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
const siteMain = "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee"
|
|
|
|
func base() time.Time { return time.Date(2026, 9, 2, 10, 0, 0, 0, time.UTC) }
|
|
|
|
// seedArrivals lays down n visits one second apart, each with a photo.
|
|
func seedArrivals(fs *fakeStore, n int) {
|
|
for i := 0; i < n; i++ {
|
|
fs.arrivals = append(fs.arrivals, Arrival{
|
|
VisitID: fmt.Sprintf("00000000-0000-0000-0000-%012d", i),
|
|
Seq: int64(i + 1),
|
|
OccurredAt: base().Add(time.Duration(i) * time.Second).Format(time.RFC3339Nano),
|
|
SiteID: siteMain,
|
|
Site: "Anna Nagar",
|
|
CameraID: "door",
|
|
VisitorID: fmt.Sprintf("11111111-0000-0000-0000-%012d", i),
|
|
Label: fmt.Sprintf("Visitor %d", i),
|
|
ImageKey: fmt.Sprintf("behavision/v2/acme/main/2026/09/02/%d.jpg", i),
|
|
})
|
|
}
|
|
}
|
|
|
|
func getPage(t *testing.T, s *Server, path, token string) ArrivalPage {
|
|
t.Helper()
|
|
rec := do(t, s, "GET", path, token, nil)
|
|
if rec.Code != http.StatusOK {
|
|
t.Fatalf("GET %s: got %d, body %s", path, rec.Code, rec.Body.String())
|
|
}
|
|
var page ArrivalPage
|
|
if err := json.Unmarshal(rec.Body.Bytes(), &page); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return page
|
|
}
|
|
|
|
// Four people walking through a door together is the case this endpoint was
|
|
// built for, and the one that used to take nine requests to render: a search
|
|
// that could not tell you who arrived, then one image call per person.
|
|
func TestFourPeopleArrivingTogetherComeBackInOneRequest(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Blob = &fakeBlob{}
|
|
seedUser(fs)
|
|
// All four in the same millisecond - a burst is exactly when this happens.
|
|
for i := 0; i < 4; i++ {
|
|
fs.arrivals = append(fs.arrivals, Arrival{
|
|
VisitID: fmt.Sprintf("00000000-0000-0000-0000-%012d", i),
|
|
Seq: int64(i + 1),
|
|
OccurredAt: base().Format(time.RFC3339Nano),
|
|
SiteID: siteMain, Site: "Anna Nagar",
|
|
VisitorID: fmt.Sprintf("11111111-0000-0000-0000-%012d", i),
|
|
ImageKey: fmt.Sprintf("k%d.jpg", i),
|
|
})
|
|
}
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
page := getPage(t, s, "/api/visits", sess.Token)
|
|
if len(page.Arrivals) != 4 {
|
|
t.Fatalf("want 4 arrivals in one response, got %d", len(page.Arrivals))
|
|
}
|
|
for i, a := range page.Arrivals {
|
|
if !a.Image.Available || a.Image.URL == "" {
|
|
t.Errorf("arrival %d has no photo link: %+v", i, a.Image)
|
|
}
|
|
}
|
|
// Identical timestamps must still produce a usable cursor, or a burst
|
|
// wedges the feed forever at the same instant.
|
|
if page.Cursor == "" {
|
|
t.Fatal("a burst at one instant produced no cursor")
|
|
}
|
|
seq, err := decodeCursor(page.Cursor)
|
|
if err != nil {
|
|
t.Fatalf("cursor from a burst is unreadable: %v", err)
|
|
}
|
|
// Against the SEEDED position, not one read back off the wire: `seq` is
|
|
// json:"-" because it counts every visit on the platform, so a client can
|
|
// no longer see it - and the cursor is the whole reason it does not need to.
|
|
if want := int64(4); seq != want {
|
|
t.Errorf("cursor should point at the LAST row of the burst, got %d want %d",
|
|
seq, want)
|
|
}
|
|
}
|
|
|
|
// The property the whole feed rests on: poll twice and you see every person
|
|
// exactly once, even when more arrive than fit in one page.
|
|
func TestCursorLosesNobodyAndRepeatsNobody(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
seedArrivals(fs, 25)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
seen := map[string]int{}
|
|
cursor := ""
|
|
for poll := 0; poll < 5; poll++ {
|
|
path := "/api/visits?limit=10"
|
|
if cursor != "" {
|
|
path += "&cursor=" + cursor
|
|
} else {
|
|
// Start from position zero so the walk covers every row rather
|
|
// than starting at the newest window.
|
|
path += "&cursor=" + encodeCursor(0) + "&limit=10"
|
|
}
|
|
page := getPage(t, s, path, sess.Token)
|
|
for _, a := range page.Arrivals {
|
|
seen[a.VisitID]++
|
|
}
|
|
cursor = page.Cursor
|
|
}
|
|
|
|
if len(seen) != 25 {
|
|
t.Fatalf("walked the feed and saw %d of 25 people", len(seen))
|
|
}
|
|
for id, n := range seen {
|
|
if n != 1 {
|
|
t.Errorf("visit %s delivered %d times, want exactly 1", id, n)
|
|
}
|
|
}
|
|
}
|
|
|
|
// An app that has just opened wants the last few arrivals, not the first few
|
|
// ever recorded - but still ascending, so its cursor handling is the same on
|
|
// the first poll as on every one after.
|
|
func TestFirstPollReturnsTheNewestWindowAscending(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
seedArrivals(fs, 25)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
page := getPage(t, s, "/api/visits?limit=5", sess.Token)
|
|
if len(page.Arrivals) != 5 {
|
|
t.Fatalf("want 5, got %d", len(page.Arrivals))
|
|
}
|
|
if got := page.Arrivals[0].VisitID; !strings.HasSuffix(got, "020") {
|
|
t.Errorf("first poll should start at the 21st of 25 rows, got %s", got)
|
|
}
|
|
for i := 1; i < len(page.Arrivals); i++ {
|
|
if page.Arrivals[i-1].OccurredAt >= page.Arrivals[i].OccurredAt {
|
|
t.Fatalf("feed is not ascending at %d", i)
|
|
}
|
|
}
|
|
}
|
|
|
|
// A quiet minute must not reset the feed. Handing back an empty cursor would
|
|
// make the next poll re-deliver the whole recent window.
|
|
func TestAnEmptyPollHandsTheCursorBack(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
seedArrivals(fs, 3)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
first := getPage(t, s, "/api/visits", sess.Token)
|
|
again := getPage(t, s, "/api/visits?cursor="+first.Cursor, sess.Token)
|
|
|
|
if len(again.Arrivals) != 0 {
|
|
t.Fatalf("nothing new arrived, got %d rows", len(again.Arrivals))
|
|
}
|
|
if again.Cursor != first.Cursor {
|
|
t.Errorf("empty poll moved the cursor: %q -> %q", first.Cursor, again.Cursor)
|
|
}
|
|
}
|
|
|
|
// The object key names a tenant's storage prefix and is the input to every
|
|
// signing call. It must be structurally incapable of reaching a client.
|
|
func TestTheObjectKeyIsNeverInTheResponse(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Blob = &fakeBlob{}
|
|
seedUser(fs)
|
|
seedArrivals(fs, 3)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
rec := do(t, s, "GET", "/api/visits", sess.Token, nil)
|
|
body := rec.Body.String()
|
|
if strings.Contains(body, `"image_key"`) || strings.Contains(body, "ImageKey") {
|
|
t.Fatalf("the response carries an object key:\n%s", body)
|
|
}
|
|
// The signed URL legitimately contains the key; what must not appear is a
|
|
// bare key in its own field.
|
|
if !strings.Contains(body, "X-Amz-Signature") {
|
|
t.Fatalf("expected signed links in the response:\n%s", body)
|
|
}
|
|
}
|
|
|
|
// Photos are off by default across the product, so "no photo" is the normal
|
|
// case and must not read as a fault. Two different absences need two different
|
|
// sentences, because a shop can fix one of them and not the other.
|
|
func TestNoPhotoIsDataNotAnError(t *testing.T) {
|
|
t.Run("images switched off for the deployment", func(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Blob = nil
|
|
seedUser(fs)
|
|
seedArrivals(fs, 1)
|
|
// No key, because that is what this deployment actually produces: the
|
|
// engine's `app.store_faces` is off, so no crop is ever captured and no
|
|
// key is ever written. A bucket key on a server with no bucket is a
|
|
// different state entirely and gets its own sentence below.
|
|
fs.arrivals[0].ImageKey = ""
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
page := getPage(t, s, "/api/visits", sess.Token)
|
|
got := page.Arrivals[0].Image
|
|
if got.Available || got.Reason == "" {
|
|
t.Fatalf("want an unavailable photo with a reason, got %+v", got)
|
|
}
|
|
if !strings.Contains(got.Reason, "not storing") {
|
|
t.Errorf("reason should say the system stores no photos, got %q", got.Reason)
|
|
}
|
|
})
|
|
|
|
// Three absences now, not two: face images may live in a bucket OR in this
|
|
// database, so "there is no bucket" stopped being a synonym for "there are
|
|
// no photos" the moment the fallback existed.
|
|
t.Run("a bucket key on a server that has lost its bucket", func(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Blob = nil
|
|
seedUser(fs)
|
|
seedArrivals(fs, 1) // seeded with an object-store key
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
page := getPage(t, s, "/api/visits", sess.Token)
|
|
got := page.Arrivals[0].Image
|
|
if got.Available {
|
|
t.Fatalf("nothing can be served without the bucket, got %+v", got)
|
|
}
|
|
// Deliberately NOT "we store no photos". The photo exists and this
|
|
// server can no longer reach it, which is a configuration fault
|
|
// somebody can fix - and reporting it as an ordinary empty record is
|
|
// how it would go unnoticed for a year.
|
|
if !strings.Contains(got.Reason, "no longer reach") {
|
|
t.Errorf("want a configuration reason, got %q", got.Reason)
|
|
}
|
|
})
|
|
|
|
t.Run("a face this server holds itself", func(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Blob = nil // no object storage anywhere
|
|
seedUser(fs)
|
|
seedArrivals(fs, 1)
|
|
fs.arrivals[0].ImageKey = "db:00000000-0000-4000-b000-000000000001"
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
page := getPage(t, s, "/api/visits", sess.Token)
|
|
got := page.Arrivals[0].Image
|
|
if !got.Available {
|
|
t.Fatalf("a stored face should be offered, got %+v", got)
|
|
}
|
|
// Auth is what tells a client this URL needs the session bearer. A
|
|
// browser <img> cannot load it and a mobile image view can, and there
|
|
// is nothing in the URL itself that says so.
|
|
if !got.Auth {
|
|
t.Error("a face held by this server must be marked as needing auth")
|
|
}
|
|
if strings.Contains(got.URL, "db:") {
|
|
t.Errorf("the storage key leaked into the URL: %q", got.URL)
|
|
}
|
|
})
|
|
|
|
t.Run("this visit simply had none", func(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Blob = &fakeBlob{}
|
|
seedUser(fs)
|
|
seedArrivals(fs, 1)
|
|
fs.arrivals[0].ImageKey = ""
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
page := getPage(t, s, "/api/visits", sess.Token)
|
|
got := page.Arrivals[0].Image
|
|
if got.Available {
|
|
t.Fatalf("there is no key, so there is no photo: %+v", got)
|
|
}
|
|
if !strings.Contains(got.Reason, "No photo") {
|
|
t.Errorf("want a per-visit reason, got %q", got.Reason)
|
|
}
|
|
})
|
|
}
|
|
|
|
// A visit with no visitor_id is a site reporting footfall without templates.
|
|
// It is a real person walking in and must appear, or the feed disagrees with
|
|
// the footfall report about how many people came.
|
|
func TestAnUnidentifiedVisitStillAppears(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Blob = &fakeBlob{}
|
|
seedUser(fs)
|
|
seedArrivals(fs, 1)
|
|
fs.arrivals[0].VisitorID = ""
|
|
fs.arrivals[0].Label = ""
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
page := getPage(t, s, "/api/visits", sess.Token)
|
|
if len(page.Arrivals) != 1 {
|
|
t.Fatalf("an anonymous visit was dropped from the feed")
|
|
}
|
|
}
|
|
|
|
// A tablet polling every two seconds would write tens of thousands of audit
|
|
// rows a day and bury the one deliberate look an investigation is after.
|
|
func TestTheFeedWritesOneAuditRowPerPageNotPerPhoto(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Blob = &fakeBlob{}
|
|
seedUser(fs)
|
|
seedArrivals(fs, 6)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
getPage(t, s, "/api/visits", sess.Token)
|
|
|
|
fs.mu.Lock()
|
|
defer fs.mu.Unlock()
|
|
var views []AuditEntry
|
|
for _, a := range fs.audits {
|
|
if strings.HasPrefix(a.Action, "image.view") {
|
|
views = append(views, a)
|
|
}
|
|
}
|
|
if len(views) != 1 {
|
|
t.Fatalf("want exactly 1 audit row for a page of 6 photos, got %d", len(views))
|
|
}
|
|
if got := views[0].Detail["count"]; got != 6 {
|
|
t.Errorf("the audit row should record how many faces were surfaced, got %v", got)
|
|
}
|
|
}
|
|
|
|
// Nothing is audited when no face was actually shown - otherwise the log fills
|
|
// with rows recording that somebody looked at nothing.
|
|
func TestAQuietPollWritesNoAuditRow(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Blob = &fakeBlob{}
|
|
seedUser(fs)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
getPage(t, s, "/api/visits", sess.Token)
|
|
|
|
fs.mu.Lock()
|
|
defer fs.mu.Unlock()
|
|
for _, a := range fs.audits {
|
|
if strings.HasPrefix(a.Action, "image.view") {
|
|
t.Fatalf("audited an image view with no images: %+v", a)
|
|
}
|
|
}
|
|
}
|
|
|
|
// The tenant comes from the session. A site_id in the query string is
|
|
// caller-controlled and is a cross-tenant read the moment it is trusted alone.
|
|
func TestTheFeedIsScopedToTheSessionsTenant(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
seedArrivals(fs, 3)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
getPage(t, s, "/api/visits?site_id="+siteMain, sess.Token)
|
|
|
|
fs.mu.Lock()
|
|
defer fs.mu.Unlock()
|
|
if fs.arrivalQ.ClientID != "client-acme" {
|
|
t.Fatalf("query ran for client %q, want the session's own", fs.arrivalQ.ClientID)
|
|
}
|
|
}
|
|
|
|
func TestTheFeedNeedsASession(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
if rec := do(t, s, "GET", "/api/visits", "", nil); rec.Code != http.StatusUnauthorized {
|
|
t.Fatalf("want 401 without a token, got %d", rec.Code)
|
|
}
|
|
}
|
|
|
|
// A malformed cursor is the caller's, not a server fault, and the message has
|
|
// to tell a client how to recover - it cannot parse the cursor to fix it.
|
|
func TestARubbishCursorIsARecoverableBadRequest(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
for _, bad := range []string{"not-base64!!", "Zm9v",
|
|
base64.RawURLEncoding.EncodeToString([]byte("v1:banana"))} {
|
|
rec := do(t, s, "GET", "/api/visits?cursor="+bad, sess.Token, nil)
|
|
if rec.Code != http.StatusBadRequest {
|
|
t.Errorf("cursor %q: got %d, want 400", bad, rec.Code)
|
|
}
|
|
if !strings.Contains(rec.Body.String(), "poll again without one") {
|
|
t.Errorf("cursor %q: message does not say how to recover: %s", bad, rec.Body.String())
|
|
}
|
|
}
|
|
}
|
|
|
|
// The id goes into $4::uuid. A malformed one is a Postgres cast error, which
|
|
// surfaces as a 500 on a value the caller supplied.
|
|
// A cursor whose body is not a number must be refused rather than reaching the
|
|
// query, and a cursor from a FUTURE encoding must fail cleanly rather than being
|
|
// misread as a position that now means something else.
|
|
func TestOnlyAWellFormedV1CursorIsAccepted(t *testing.T) {
|
|
for _, bad := range []string{
|
|
base64.RawURLEncoding.EncodeToString([]byte("v1:not-a-number")),
|
|
base64.RawURLEncoding.EncodeToString([]byte("v1:-5")),
|
|
base64.RawURLEncoding.EncodeToString([]byte("v2:12")),
|
|
base64.RawURLEncoding.EncodeToString([]byte("12")),
|
|
base64.RawURLEncoding.EncodeToString([]byte("v1:'; DROP TABLE visits; --")),
|
|
} {
|
|
if _, err := decodeCursor(bad); err == nil {
|
|
t.Errorf("accepted a malformed cursor: %q", bad)
|
|
}
|
|
}
|
|
if seq, err := decodeCursor(encodeCursor(41)); err != nil || seq != 41 {
|
|
t.Fatalf("a cursor did not round-trip: %d %v", seq, err)
|
|
}
|
|
}
|
|
|
|
func TestASiteFilterMustBeASiteID(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
rec := do(t, s, "GET", "/api/visits?site_id=main", sess.Token, nil)
|
|
if rec.Code != http.StatusBadRequest {
|
|
t.Fatalf("want 400 for a non-uuid site_id, got %d", rec.Code)
|
|
}
|
|
}
|
|
|
|
func TestLimitIsCappedRatherThanRejected(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
// A client coming back from a tunnel asking for everything gets a big page
|
|
// and a fresh cursor, not an error it cannot recover from.
|
|
getPage(t, s, "/api/visits?limit=100000", sess.Token)
|
|
|
|
fs.mu.Lock()
|
|
defer fs.mu.Unlock()
|
|
if fs.arrivalQ.Limit != maxArrivals {
|
|
t.Fatalf("limit %d, want it capped at %d", fs.arrivalQ.Limit, maxArrivals)
|
|
}
|
|
}
|
|
|
|
func TestADatabaseFailureIsAServerErrorNotAnEmptyFeed(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
fs.arrivalsErr = errors.New("connection refused")
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
rec := do(t, s, "GET", "/api/visits", sess.Token, nil)
|
|
if rec.Code != http.StatusInternalServerError {
|
|
// An empty 200 here would tell a shop nobody came in.
|
|
t.Fatalf("want 500, got %d", rec.Code)
|
|
}
|
|
if strings.Contains(rec.Body.String(), "connection refused") {
|
|
t.Error("the database error leaked to the client")
|
|
}
|
|
}
|
|
|
|
func TestArrivalsIsAlwaysAListNeverNull(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
rec := do(t, s, "GET", "/api/visits", sess.Token, nil)
|
|
if !strings.Contains(rec.Body.String(), `"arrivals":[]`) {
|
|
t.Fatalf("a quiet feed must be an empty list: %s", rec.Body.String())
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------- streaming
|
|
|
|
func TestTheStreamSendsWhatIsAlreadyThereThenRespondsToTheDoorbell(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Blob = &fakeBlob{}
|
|
s.Hub = NewHub()
|
|
seedUser(fs)
|
|
seedArrivals(fs, 2)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
defer cancel()
|
|
req := httptest.NewRequest("GET", "/api/visits/stream", nil).WithContext(ctx)
|
|
req.Header.Set("Authorization", "Bearer "+sess.Token)
|
|
rec := newStreamRecorder()
|
|
|
|
done := make(chan struct{})
|
|
go func() { defer close(done); s.Routes().ServeHTTP(rec, req) }()
|
|
|
|
// The first frame is the catch-up: a client connecting after people have
|
|
// walked in must not be blind until the next person arrives.
|
|
waitFor(t, rec, "event: arrivals")
|
|
if h := rec.Header().Get("Content-Type"); h != "text/event-stream" {
|
|
t.Errorf("Content-Type %q", h)
|
|
}
|
|
if h := rec.Header().Get("X-Accel-Buffering"); h != "no" {
|
|
t.Error("without this the proxy buffers the stream into a single response at the end")
|
|
}
|
|
if !strings.Contains(rec.body(), "id: ") {
|
|
t.Error("no SSE id, so a reconnect cannot resume via Last-Event-ID")
|
|
}
|
|
|
|
before := len(rec.body())
|
|
fs.mu.Lock()
|
|
fs.arrivals = append(fs.arrivals, Arrival{
|
|
VisitID: "00000000-0000-0000-0000-000000000099", Seq: 99,
|
|
OccurredAt: base().Add(time.Hour).Format(time.RFC3339Nano),
|
|
SiteID: siteMain, VisitorID: "11111111-0000-0000-0000-000000000099",
|
|
})
|
|
fs.mu.Unlock()
|
|
s.Hub.Notify("client-acme")
|
|
|
|
waitForGrowth(t, rec, before)
|
|
cancel()
|
|
<-done
|
|
}
|
|
|
|
// One process serves every tenant. A doorbell for another shop must not make
|
|
// this connection query, let alone emit anything.
|
|
func TestAnotherTenantsDoorbellIsIgnored(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
s.Hub = NewHub()
|
|
seedUser(fs)
|
|
sess := login(t, s, "manager@acme.com", "correct horse battery")
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
defer cancel()
|
|
req := httptest.NewRequest("GET", "/api/visits/stream", nil).WithContext(ctx)
|
|
req.Header.Set("Authorization", "Bearer "+sess.Token)
|
|
rec := newStreamRecorder()
|
|
done := make(chan struct{})
|
|
go func() { defer close(done); s.Routes().ServeHTTP(rec, req) }()
|
|
|
|
waitForSubscriber(t, s.Hub)
|
|
fs.mu.Lock()
|
|
fs.arrivals = append(fs.arrivals, Arrival{
|
|
VisitID: "00000000-0000-0000-0000-000000000001", Seq: 1,
|
|
OccurredAt: base().Format(time.RFC3339Nano), SiteID: siteMain,
|
|
})
|
|
queriesBefore := fs.arrivalCalls
|
|
fs.mu.Unlock()
|
|
|
|
s.Hub.Notify("client-someone-else")
|
|
time.Sleep(50 * time.Millisecond)
|
|
|
|
fs.mu.Lock()
|
|
after := fs.arrivalCalls
|
|
fs.mu.Unlock()
|
|
if after != queriesBefore {
|
|
t.Fatalf("another tenant's doorbell caused %d queries", after-queriesBefore)
|
|
}
|
|
cancel()
|
|
<-done
|
|
}
|
|
|
|
func TestTheStreamNeedsASession(t *testing.T) {
|
|
s, fs := newServer(t)
|
|
seedUser(fs)
|
|
if rec := do(t, s, "GET", "/api/visits/stream", "", nil); rec.Code != http.StatusUnauthorized {
|
|
t.Fatalf("want 401, got %d", rec.Code)
|
|
}
|
|
}
|
|
|
|
// A subscriber that is not released leaks a goroutine and a channel for the
|
|
// life of the process, on an endpoint mobile clients reconnect to all day.
|
|
func TestReleasingASubscriberRemovesIt(t *testing.T) {
|
|
h := NewHub()
|
|
_, release := h.Subscribe()
|
|
if h.Subscribers() != 1 {
|
|
t.Fatalf("want 1 subscriber, got %d", h.Subscribers())
|
|
}
|
|
release()
|
|
if h.Subscribers() != 0 {
|
|
t.Fatalf("subscriber leaked: %d still registered", h.Subscribers())
|
|
}
|
|
release() // must be safe twice: defer plus an early return is normal
|
|
}
|
|
|
|
// A doorbell is idempotent, so a busy listener that misses one loses nothing -
|
|
// but Notify must never block waiting for it, or one slow subscriber stalls
|
|
// ingest for the whole estate.
|
|
func TestNotifyNeverBlocksOnASlowSubscriber(t *testing.T) {
|
|
h := NewHub()
|
|
ch, release := h.Subscribe()
|
|
defer release()
|
|
|
|
done := make(chan struct{})
|
|
go func() {
|
|
defer close(done)
|
|
for i := 0; i < 1000; i++ {
|
|
h.Notify("client-acme")
|
|
}
|
|
}()
|
|
select {
|
|
case <-done:
|
|
case <-time.After(2 * time.Second):
|
|
t.Fatal("Notify blocked - ingest would stall behind a slow reader")
|
|
}
|
|
if len(ch) != 1 {
|
|
t.Errorf("want the single pending doorbell, got %d", len(ch))
|
|
}
|
|
}
|
|
|
|
func TestNotifyOnANilHubIsSafe(t *testing.T) {
|
|
var h *Hub
|
|
h.Notify("client-acme") // a server assembled without one must still ingest
|
|
if h.Subscribers() != 0 {
|
|
t.Fatal("nil hub reported subscribers")
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------- helpers
|
|
|
|
// streamRecorder is a ResponseWriter a test can read WHILE the handler is
|
|
// still writing to it. httptest.ResponseRecorder cannot be: the handler runs on
|
|
// its own goroutine for the life of the stream, so every Body.String() from the
|
|
// test is a data race that -race turns into a failure and, without it, into an
|
|
// occasional mystery.
|
|
//
|
|
// It implements Flusher because the handler refuses to stream without one - a
|
|
// recorder that silently lacked it would make these tests exercise the error
|
|
// path while appearing to pass.
|
|
type streamRecorder struct {
|
|
mu sync.Mutex
|
|
buf strings.Builder
|
|
hdr http.Header
|
|
code int
|
|
pushed chan struct{}
|
|
}
|
|
|
|
func newStreamRecorder() *streamRecorder {
|
|
return &streamRecorder{hdr: http.Header{}, code: 200,
|
|
pushed: make(chan struct{}, 64)}
|
|
}
|
|
|
|
func (r *streamRecorder) Header() http.Header { return r.hdr }
|
|
|
|
func (r *streamRecorder) Write(b []byte) (int, error) {
|
|
r.mu.Lock()
|
|
n, err := r.buf.Write(b)
|
|
r.mu.Unlock()
|
|
return n, err
|
|
}
|
|
|
|
func (r *streamRecorder) WriteHeader(code int) { r.code = code }
|
|
|
|
func (r *streamRecorder) Flush() {
|
|
// Signals the test that a frame is complete, so it can wait on an event
|
|
// rather than on a sleep long enough to hide a real stall.
|
|
select {
|
|
case r.pushed <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
func (r *streamRecorder) body() string {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
return r.buf.String()
|
|
}
|
|
|
|
func waitFor(t *testing.T, rec *streamRecorder, want string) {
|
|
t.Helper()
|
|
deadline := time.Now().Add(2 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if strings.Contains(rec.body(), want) {
|
|
return
|
|
}
|
|
time.Sleep(5 * time.Millisecond)
|
|
}
|
|
t.Fatalf("never saw %q in the stream:\n%s", want, rec.body())
|
|
}
|
|
|
|
func waitForGrowth(t *testing.T, rec *streamRecorder, was int) {
|
|
t.Helper()
|
|
deadline := time.Now().Add(2 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if len(rec.body()) > was {
|
|
return
|
|
}
|
|
time.Sleep(5 * time.Millisecond)
|
|
}
|
|
t.Fatal("the doorbell did not push a new arrival within 2s")
|
|
}
|
|
|
|
func waitForSubscriber(t *testing.T, h *Hub) {
|
|
t.Helper()
|
|
deadline := time.Now().Add(2 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
if h.Subscribers() > 0 {
|
|
return
|
|
}
|
|
time.Sleep(5 * time.Millisecond)
|
|
}
|
|
t.Fatal("stream never subscribed to the hub")
|
|
}
|