package store import ( "context" "errors" "fmt" "strings" "time" "github.com/jackc/pgx/v5" "github.com/loyaly/behavision-server/internal/api" "github.com/loyaly/behavision-server/internal/auth" ) // ---------------------------------------------------------------- identity func (s *Store) UserByEmail(ctx context.Context, email string) (api.UserRecord, error) { var u api.UserRecord err := s.pool.QueryRow(ctx, ` SELECT u.id::text, COALESCE(u.client_id::text, ''), COALESCE(c.name, ''), u.email, u.full_name, u.role, u.active, u.password_hash FROM app_users u LEFT JOIN clients c ON c.id = u.client_id WHERE lower(u.email) = $1`, email). Scan(&u.ID, &u.ClientID, &u.ClientName, &u.Email, &u.FullName, &u.Role, &u.Active, &u.PasswordHash) if errors.Is(err, pgx.ErrNoRows) { // Not an error. The handler must still spend the same time verifying a // password, so "no such user" has to come back as data rather than as a // short-circuit. return api.UserRecord{Found: false}, nil } if err != nil { return api.UserRecord{}, err } // A user whose client has been deactivated must not be able to sign in and // read that client's customers. if u.ClientID != "" { var active bool if err := s.pool.QueryRow(ctx, `SELECT active FROM clients WHERE id = $1`, u.ClientID). Scan(&active); err != nil { return api.UserRecord{}, err } u.Active = u.Active && active } u.Found = true return u, nil } func (s *Store) TouchUserLogin(ctx context.Context, userID string) error { _, err := s.pool.Exec(ctx, `UPDATE app_users SET last_login_at = now() WHERE id = $1`, userID) return err } func (s *Store) CreateSession(ctx context.Context, n api.NewSession) error { _, err := s.pool.Exec(ctx, ` INSERT INTO sessions (user_id, client_id, access_hash, refresh_hash, access_expires_at, refresh_expires_at, device) VALUES ($1, NULLIF($2, '')::uuid, $3, $4, $5, $6, $7)`, n.UserID, n.ClientID, n.AccessHash, n.RefreshHash, n.AccessExpiry, n.RefreshExp, n.Device) return err } // sessionQuery is shared by the access and refresh lookups so the two can // never disagree about what makes a session valid. const sessionQuery = ` SELECT s.id::text, s.user_id::text, COALESCE(s.client_id::text, ''), COALESCE(c.name, ''), u.email, u.full_name, u.role, %s FROM sessions s JOIN app_users u ON u.id = s.user_id AND u.active LEFT JOIN clients c ON c.id = s.client_id WHERE s.%s = $1 AND s.revoked_at IS NULL` func (s *Store) SessionByAccess(ctx context.Context, hash []byte) (auth.Principal, time.Time, error) { // last_used_at is refreshed at most every five minutes. Writing it on every // authenticated request would turn a read-only API call into a row update // and a WAL record, for a column nothing needs to the second. if _, err := s.pool.Exec(ctx, ` UPDATE sessions SET last_used_at = now() WHERE access_hash = $1 AND revoked_at IS NULL AND (last_used_at IS NULL OR last_used_at < now() - interval '5 minutes')`, hash); err != nil { // Bookkeeping. Refusing the request because a timestamp would not // update would log everybody out over nothing. _ = err } return s.session(ctx, hash, "access_expires_at", "access_hash") } func (s *Store) SessionByRefresh(ctx context.Context, hash []byte) (auth.Principal, time.Time, error) { return s.session(ctx, hash, "refresh_expires_at", "refresh_hash") } func (s *Store) session(ctx context.Context, hash []byte, expiryCol, hashCol string) ( auth.Principal, time.Time, error) { var p auth.Principal var expires time.Time err := s.pool.QueryRow(ctx, fmt.Sprintf(sessionQuery, expiryCol, hashCol), hash). Scan(&p.SessionID, &p.UserID, &p.ClientID, &p.ClientName, &p.Email, &p.FullName, &p.Role, &expires) if errors.Is(err, pgx.ErrNoRows) { return auth.Principal{}, time.Time{}, auth.ErrNoSession } return p, expires, err } // RotateSession replaces the tokens on an existing row rather than inserting a // new one. The old refresh token stops working the moment this commits, which // is the point: a token copied off a resold shop PC must not keep working // alongside the real one. func (s *Store) RotateSession(ctx context.Context, sessionID string, n api.NewSession) error { tag, err := s.pool.Exec(ctx, ` UPDATE sessions SET access_hash = $2, refresh_hash = $3, access_expires_at = $4, refresh_expires_at = $5, last_used_at = now(), device = COALESCE(NULLIF($6, ''), device) WHERE id = $1 AND revoked_at IS NULL`, sessionID, n.AccessHash, n.RefreshHash, n.AccessExpiry, n.RefreshExp, n.Device) if err != nil { return err } if tag.RowsAffected() == 0 { return auth.ErrNoSession } return nil } func (s *Store) RevokeSession(ctx context.Context, sessionID string) error { _, err := s.pool.Exec(ctx, `UPDATE sessions SET revoked_at = now() WHERE id = $1 AND revoked_at IS NULL`, sessionID) return err } // ---------------------------------------------------------------- reports func nullUUID(s string) any { if strings.TrimSpace(s) == "" { return nil } return s } // Footfall buckets visits in the requested timezone. // // Two things here are easy to get wrong and expensive to notice: // // - "New" means first-ever, computed over all of time, not first-in-window. // Otherwise every report re-labels your regulars as new customers the // moment the window starts after their last visit. // - A visit with no visitor_id (a site sending counts without templates) is // real footfall but an unknown person. It is counted in `visitors` and in // neither `new` nor `returning`, so those two may sum to less than the // total. Guessing either way would put a number in a marketing report that // nothing supports. func (s *Store) Footfall(ctx context.Context, q api.ReportQuery) ( []api.FootfallPoint, api.Totals, error) { site := nullUUID(q.SiteID) rows, err := s.pool.Query(ctx, ` WITH scoped AS ( SELECT v.visitor_id, v.occurred_at FROM visits v WHERE v.client_id = $1 AND v.occurred_at >= $2 AND v.occurred_at < $3 AND ($4::uuid IS NULL OR v.site_id = $4::uuid) ), firsts AS ( SELECT v.visitor_id, min(v.occurred_at) AS first_at FROM visits v WHERE v.client_id = $1 AND v.visitor_id IS NOT NULL AND ($4::uuid IS NULL OR v.site_id = $4::uuid) GROUP BY v.visitor_id ) SELECT date_trunc($5, s.occurred_at AT TIME ZONE $6) AS bucket, count(DISTINCT s.visitor_id) AS identified, count(*) FILTER (WHERE s.visitor_id IS NULL) AS anonymous, count(DISTINCT s.visitor_id) FILTER ( WHERE date_trunc($5, f.first_at AT TIME ZONE $6) = date_trunc($5, s.occurred_at AT TIME ZONE $6)) AS newcomers FROM scoped s LEFT JOIN firsts f ON f.visitor_id = s.visitor_id GROUP BY 1 ORDER BY 1`, q.ClientID, q.From, q.To, site, q.Bucket, q.Timezone) if err != nil { return nil, api.Totals{}, fmt.Errorf("footfall buckets: %w", err) } defer rows.Close() var points []api.FootfallPoint for rows.Next() { var t time.Time var identified, anonymous, newcomers int if err := rows.Scan(&t, &identified, &anonymous, &newcomers); err != nil { return nil, api.Totals{}, err } points = append(points, api.FootfallPoint{ // Local wall time, with no offset, because the label belongs to the // timezone named alongside it in the report. Stamping it with Z // would say 09:00 UTC when the shop means 09:00 in Chennai. Bucket: t.Format("2006-01-02T15:04:05"), Visitors: identified + anonymous, New: newcomers, Returning: identified - newcomers, }) } if err := rows.Err(); err != nil { return nil, api.Totals{}, err } if points == nil { points = []api.FootfallPoint{} } var totals api.Totals var identified, anonymous int if err := s.pool.QueryRow(ctx, ` SELECT count(DISTINCT v.visitor_id), count(*) FILTER (WHERE v.visitor_id IS NULL), count(*) FROM visits v WHERE v.client_id = $1 AND v.occurred_at >= $2 AND v.occurred_at < $3 AND ($4::uuid IS NULL OR v.site_id = $4::uuid)`, q.ClientID, q.From, q.To, site). Scan(&identified, &anonymous, &totals.Visits); err != nil { return nil, api.Totals{}, fmt.Errorf("footfall totals: %w", err) } totals.UniqueVisitors = identified + anonymous // Worst site, not the average. One badly placed camera is a hole in this // report, and averaging it against three good ones hides the only site // anyone needs to do something about. err = s.pool.QueryRow(ctx, ` SELECT a.fraction_below_gate, si.name FROM agents a JOIN sites si ON si.id = a.site_id WHERE a.client_id = $1 AND a.fraction_below_gate IS NOT NULL AND ($2::uuid IS NULL OR a.site_id = $2::uuid) ORDER BY a.fraction_below_gate DESC LIMIT 1`, q.ClientID, site). Scan(&totals.FractionBelowGate, &totals.WorstSite) if err != nil && !errors.Is(err, pgx.ErrNoRows) { return nil, api.Totals{}, fmt.Errorf("gate fraction: %w", err) } return points, totals, nil } // Conversion answers "how many of the people who walked in bought something". // // Revenue is summed for ONE currency - whichever accounts for the most of it. // Adding rupees to dollars produces a number that looks like money and is not, // and this figure is the one a customer judges the product by. func (s *Store) Conversion(ctx context.Context, q api.ReportQuery) (api.SalesReport, error) { site := nullUUID(q.SiteID) var rep api.SalesReport var identified, anonymous int if err := s.pool.QueryRow(ctx, ` SELECT count(DISTINCT v.visitor_id), count(*) FILTER (WHERE v.visitor_id IS NULL) FROM visits v WHERE v.client_id = $1 AND v.occurred_at >= $2 AND v.occurred_at < $3 AND ($4::uuid IS NULL OR v.site_id = $4::uuid)`, q.ClientID, q.From, q.To, site).Scan(&identified, &anonymous); err != nil { return rep, fmt.Errorf("conversion visitors: %w", err) } rep.Visitors = identified + anonymous var baskets int if err := s.pool.QueryRow(ctx, ` WITH scoped AS ( SELECT p.visitor_id, p.amount, p.currency FROM purchases p WHERE p.client_id = $1 AND p.occurred_at >= $2 AND p.occurred_at < $3 AND ($4::uuid IS NULL OR p.site_id = $4::uuid) ), dominant AS ( SELECT currency FROM scoped GROUP BY currency ORDER BY sum(amount) DESC LIMIT 1 ) SELECT COALESCE((SELECT currency FROM dominant), 'INR'), count(DISTINCT s.visitor_id), count(*), COALESCE(sum(s.amount), 0)::float8 FROM scoped s WHERE s.currency = COALESCE((SELECT currency FROM dominant), 'INR')`, q.ClientID, q.From, q.To, site). Scan(&rep.Currency, &rep.Purchasers, &baskets, &rep.Revenue); err != nil { return rep, fmt.Errorf("conversion purchases: %w", err) } if rep.Visitors > 0 { rep.Conversion = float64(rep.Purchasers) / float64(rep.Visitors) } if baskets > 0 { // Per basket, not per purchaser: a customer who bought twice in the // window had two baskets, and averaging over people would overstate // what a single transaction is worth. rep.AvgBasket = rep.Revenue / float64(baskets) } return rep, nil } func (s *Store) SiteHealth(ctx context.Context, clientID string) ([]api.SiteHealth, error) { rows, err := s.pool.Query(ctx, ` SELECT si.id::text, si.slug, si.name, si.timezone, a.last_heartbeat_at, a.last_event_at, COALESCE(a.recognition_model, ''), COALESCE(a.agent_version, ''), COALESCE(a.cameras_up, 0), COALESCE(a.cameras_total, 0), a.fraction_below_gate, COALESCE(a.spool_queued, 0), COALESCE(a.spool_dropped, 0) FROM sites si LEFT JOIN agents a ON a.site_id = si.id WHERE si.client_id = $1 AND si.active ORDER BY si.name`, clientID) if err != nil { return nil, err } defer rows.Close() var out []api.SiteHealth for rows.Next() { var h api.SiteHealth var beat, event *time.Time var gate *float64 if err := rows.Scan(&h.SiteID, &h.Slug, &h.Name, &h.Timezone, &beat, &event, &h.RecognitionModel, &h.AgentVersion, &h.CamerasUp, &h.CamerasTotal, &gate, &h.Queued, &h.Dropped); err != nil { return nil, err } if beat != nil { h.LastHeartbeatAt = beat.UTC().Format(time.RFC3339) // Three missed beats. One missed beat is a dropped packet; three is // a site that has actually gone away, and calling that out too // eagerly trains people to ignore the indicator. h.Online = time.Since(*beat) < 3*time.Minute } if event != nil { h.LastEventAt = event.UTC().Format(time.RFC3339) } if gate != nil { h.FractionBelowGate = *gate } out = append(out, h) } return out, rows.Err() } // Compile-time proof that the store satisfies what the API asks for. Without // it a missing method is only discovered when main.go is wired up, which is the // one file least covered by tests. var _ api.Store = (*Store)(nil)