package store import ( "context" "time" "github.com/loyaly/behavision-server/internal/api" ) // arrivalColumns is shared by both directions of the query below so the two // cannot drift apart - a column present in one and missing from the other would // mean the first poll of a feed and every poll after it returned different // shapes, which is the kind of bug that only shows up under load. const arrivalColumns = ` vi.id::text, vi.seq, COALESCE(vi.number, 0), vi.occurred_at, vi.site_id::text, si.name, si.slug, vi.camera_id, vi.is_new_visitor, vi.similarity, vi.quality, vi.attributes, vi.image_key, COALESCE(vi.visitor_id::text, ''), COALESCE(vs.number, 0), COALESCE(vs.label, ''), COALESCE(p.full_name, '')` const arrivalFrom = ` FROM visits vi JOIN sites si ON si.id = vi.site_id -- LEFT, not INNER, three times over. A visit with no visitor_id is a site -- reporting footfall without templates; an erased customer has their -- visitor row flagged deleted. Both are real arrivals and an inner join -- would silently drop them, making the feed disagree with the footfall -- report about how many people came in. LEFT JOIN visitors vs ON vs.id = vi.visitor_id AND vs.deleted_at IS NULL LEFT JOIN visitor_profiles p ON p.visitor_id = vi.visitor_id AND p.client_id = vi.client_id WHERE vi.client_id = $1 AND ($2 = '' OR vi.site_id = $2::uuid)` // Arrivals reads a window of the live feed, oldest first. // // Ordered by `seq` - the server-assigned position - and never by occurred_at. // That is the whole correctness argument for this endpoint and it is not // obvious, so: // // - occurred_at is the CAMERA's clock. Four people through one door share it // to the microsecond, so it cannot order them; and a site that was offline // for a day floods in carrying yesterday's timestamps, which a reader whose // cursor has passed them would skip entirely. // - A (occurred_at, id) tie-break does not save it either, because id is a // random uuid: a row that COMMITS after the reader moved its cursor but // carries a lower uuid sorts behind that cursor and is never delivered. // Measured live before this was fixed - four simultaneous visits, two // delivered, and nothing downstream able to tell. // // So the feed is ordered by when the server LEARNED of a visit. Each row still // carries occurred_at for display; seq is only ever a position. // // This depends on visits being inserted one at a time, which the MQTT consumer // guarantees with SetOrderMatters(true) - a single ordered handler goroutine, // so seq order is commit order. Running two server instances against one // database would break that assumption, and the fix then is a commit-ordered // cursor, not a bigger sequence. // // Keyset, never OFFSET: rows arrive into this table continuously, so an offset // shifts under the caller between polls and a feed built on it both repeats and // skips people. func (s *Store) Arrivals(ctx context.Context, q api.ArrivalQuery) ([]api.Arrival, error) { var sql string var args []any switch { case q.AfterSeq != nil: sql = `SELECT ` + arrivalColumns + arrivalFrom + ` AND vi.seq > $3 ORDER BY vi.seq ASC LIMIT $4` args = []any{q.ClientID, q.SiteID, *q.AfterSeq, q.Limit} case q.Since != nil: // "Everything I have not been told about since this instant." Resolved // against received_at, not occurred_at, so it means the same thing as // the cursor it turns into on the next poll - a caller must not get a // different feed depending on which of the two it started with. sql = `SELECT ` + arrivalColumns + arrivalFrom + ` AND vi.seq > COALESCE( (SELECT max(v2.seq) FROM visits v2 WHERE v2.client_id = $1 AND v2.received_at < $3), 0) ORDER BY vi.seq ASC LIMIT $4` args = []any{q.ClientID, q.SiteID, *q.Since, q.Limit} default: // No cursor: an app that has just opened. It wants the last few // arrivals, not the first few ever recorded, so take the newest rows // and reverse them - the response is still ascending, so the caller's // cursor handling is identical on the first poll and every one after. sql = `SELECT * FROM ( SELECT ` + arrivalColumns + arrivalFrom + ` ORDER BY vi.seq DESC LIMIT $3 ) t ORDER BY t.seq ASC` args = []any{q.ClientID, q.SiteID, q.Limit} } rows, err := s.pool.Query(ctx, sql, args...) if err != nil { return nil, err } defer rows.Close() var out []api.Arrival for rows.Next() { var a api.Arrival var at time.Time var sim, qual *float64 var imageKey string var number, visitNumber int64 if err := rows.Scan(&a.VisitID, &a.Seq, &visitNumber, &at, &a.SiteID, &a.Site, &a.SiteSlug, &a.CameraID, &a.IsNew, &sim, &qual, &a.Attributes, &imageKey, &a.VisitorID, &number, &a.Label, &a.Name); err != nil { return nil, err } a.VisitorRef = api.VisitorRef(number) a.VisitRef = api.VisitRef(visitNumber) a.OccurredAt = at.UTC().Format(time.RFC3339Nano) if sim != nil { a.Similarity = *sim } if qual != nil { a.Quality = *qual } // The store never presigns. It has no bucket and no idea whether this // caller is allowed to look, and a query that mints credentials is one // refactor away from doing it on a path that never checked. ImageKey // is json:"-", so a handler that forgets to swap it leaks nothing. a.ImageKey = imageKey out = append(out, a) } return out, rows.Err() }