package store import ( "context" "errors" "fmt" "time" "github.com/jackc/pgx/v5" ) // ErrNoSnapshot means this camera has no stored picture. It is an ordinary // state - a camera added a minute ago has none - so callers report it as // absence rather than as a failure. var ErrNoSnapshot = errors.New("no snapshot for this camera") // PutCameraSnapshot stores the latest frame from one of a site's cameras. // // The camera is resolved by (site_id, camera_id) IN THE INSERT, so an agent // physically cannot store a picture against another site's camera even if it // sends one - the same rule as every other agent-authenticated write here. // `camera_id` is what the ENGINE knows the camera by, because that is the only // name the shop PC has. func (s *Store) PutCameraSnapshot(ctx context.Context, clientID, siteID, cameraID string, jpeg []byte) error { tx, err := s.pool.Begin(ctx) if err != nil { return err } defer func() { _ = tx.Rollback(context.WithoutCancel(ctx)) }() var id string err = tx.QueryRow(ctx, ` SELECT id::text FROM site_cameras WHERE site_id = $1::uuid AND camera_id = $2 AND deleted_at IS NULL`, siteID, cameraID).Scan(&id) if errors.Is(err, pgx.ErrNoRows) { // Head office has not been told about this camera yet, or it was // removed. Neither is an error the agent can act on: the next sync // adopts it and the snapshot after that lands. return ErrNoSnapshot } if err != nil { return err } if _, err := tx.Exec(ctx, ` INSERT INTO camera_snapshots (camera_id, client_id, site_id, image, bytes, captured_at) VALUES ($1::uuid, $2::uuid, $3::uuid, $4, $5, now()) ON CONFLICT (camera_id) DO UPDATE SET image = EXCLUDED.image, bytes = EXCLUDED.bytes, captured_at = EXCLUDED.captured_at`, id, clientID, siteID, jpeg, len(jpeg)); err != nil { return fmt.Errorf("store snapshot: %w", err) } // snapshot_at is what tells the camera list a picture exists at all, and it // is written in the SAME transaction as the bytes. Set apart, a camera // could advertise a picture that is not there - which renders as a broken // image on the one screen whose job is to show the camera. if _, err := tx.Exec(ctx, ` UPDATE site_cameras SET snapshot_at = now() WHERE id = $1::uuid`, id); err != nil { return err } return tx.Commit(ctx) } // CameraSnapshot returns a camera's stored picture, scoped to the tenant. func (s *Store) CameraSnapshot(ctx context.Context, clientID, cameraID string) ( []byte, time.Time, error) { var img []byte var at time.Time err := s.pool.QueryRow(ctx, ` SELECT image, captured_at FROM camera_snapshots WHERE camera_id = $1::uuid AND client_id = $2::uuid`, cameraID, clientID).Scan(&img, &at) if errors.Is(err, pgx.ErrNoRows) { return nil, time.Time{}, ErrNoSnapshot } return img, at, err } // CameraRef resolves one of a tenant's cameras to its site and the name the // engine knows it by. // // Its job is to prove ownership before a live stream starts. Everything after // that point is keyed on a camera id, and a hub relaying frames does not know // whose camera it is holding — so this is the only place that can decide. func (s *Store) CameraRef(ctx context.Context, clientID, cameraID string) (string, string, error) { return s.cameraRef(ctx, `client_id = $2::uuid`, cameraID, clientID) } // CameraRefBySite is the same question asked by an agent, which is // authenticated for a site rather than a tenant. func (s *Store) CameraRefBySite(ctx context.Context, siteID, cameraID string) (string, string, error) { return s.cameraRef(ctx, `site_id = $2::uuid`, cameraID, siteID) } func (s *Store) cameraRef(ctx context.Context, scope, cameraID, owner string) (string, string, error) { var siteID, engineID string err := s.pool.QueryRow(ctx, ` SELECT site_id::text, camera_id FROM site_cameras WHERE id = $1::uuid AND `+scope+` AND deleted_at IS NULL`, cameraID, owner).Scan(&siteID, &engineID) if errors.Is(err, pgx.ErrNoRows) { return "", "", ErrNoSnapshot } return siteID, engineID, err } // SiteCameraIDs lists a site's camera uuids, for the agent's live poll. func (s *Store) SiteCameraIDs(ctx context.Context, siteID string) ([]string, error) { rows, err := s.pool.Query(ctx, ` SELECT id::text FROM site_cameras WHERE site_id = $1::uuid AND deleted_at IS NULL AND enabled`, siteID) if err != nil { return nil, err } defer rows.Close() var out []string for rows.Next() { var id string if err := rows.Scan(&id); err != nil { return nil, err } out = append(out, id) } return out, rows.Err() }