changh95's picture
tt-model push diffusion-planner-p150 (container)
4d9b003 verified
Raw History Blame Contribute Delete
14.3 kB
# SPDX-License-Identifier: Apache-2.0
"""C11 point clouds: Autoware's multi-sweep densification state machine, input hygiene and the seeded permutation.
**Densification** (``PointCloudDensification`` + ``VoxelGenerator::generateSweepPoints`` of
``autoware_lidar_centerpoint``; the same code shape is in lidar_transfusion, bevfusion and
image_projection_based_fusion: S:centerpoint:164-176; S:transfusion:99-103; S:pointpainting:216-221):
- a cache holds the current cloud and ``num_past_frames`` previous ones, newest first (``push_front`` /
``pop_back``, ``pointcloud_densification.cpp:79-95``);
- with a cache size > 1 every enqueue needs the TF ``T_sensor_from_world`` at the cloud's stamp; when it is missing
the frame is **skipped and not cached** (``:66-71``). The node keeps the world->current transform and each entry's
past->world transform as ``Eigen::Affine3f`` (float32, inverted in float32: ``:47-51,82-86``);
- at the end of each frame every cached cloud is re-transformed with ``world2current(now) @ past2world(entry)``
(float32 product) by the float32 CUDA kernel, and gets ``time_lag = float(t_now - t_entry)`` in seconds
(``voxel_generator.cpp:56-93``). There is no maximum-age check: after dropped frames the past sweep can be old;
- points are written current sweep first; a sweep that would exceed ``cloud_capacity`` is dropped **with every older
sweep** (``voxel_generator.cpp:70-76``, a ``break``);
- cold start: the first frame has a single sweep (all ``time_lag = 0``); ``num_past_frames = 0`` needs no TF at all.
:class:`SweepDensifier` is that state machine for one stream, :class:`StreamDensifiers` a per-``stream.id`` dict with
LRU eviction (PLAN.md D16: host sweep caches are per id), and :func:`densify_sweeps` the stateless variant for a
client that sends its past sweeps with ``T_current_from_sweep`` and ``time_lag_s`` itself.
**Hygiene.** Autoware's range test is written as a rejection (``x < min || x >= max || ...``), so a NaN coordinate
passes it (S:centerpoint:179); the ports drop non-finite rows explicitly (a documented deviation; ``ttaw.io.PointCloud``
already does it at decode time). Organized clouds mark missing returns as (0, 0, 0): :func:`nonzero_rows` drops them,
as the research references did (``centerpoint_ref.py``), before any transform.
**Seeded permutation** (S:centerpoint:381-389; PLAN.md D1): Autoware reads its points through a fixed
``std::shuffle`` permutation at a time-seeded random offset, so which <= 32 points an over-full pillar keeps is random.
The ports keep the statistics but are deterministic: ``numpy.random.default_rng(seed).permutation(n)`` before the
first-K-points rule (:func:`seeded_permutation`).
numpy only; no side effects on import.
"""
from __future__ import annotations
import collections
from dataclasses import dataclass, field
from typing import Any, Dict, List, Optional, Sequence, Tuple
import numpy as np
from .geometry import compose_f32, invert_affine_f32, invert_rigid, transform_points_f32
__all__ = [
"DEFAULT_CLOUD_CAPACITY",
"SweepDensifier",
"StreamDensifiers",
"DensifyInfo",
"densify_sweeps",
"finite_rows",
"nonzero_rows",
"range_mask",
"seeded_permutation",
]
DEFAULT_CLOUD_CAPACITY = 2_000_000 # centerpoint_common.param.yaml:3 (cloud_capacity)
@dataclass
class DensifyInfo:
"""What one densification produced: per used sweep its point count and time lag, plus the dropped sweeps."""
num_points: int = 0
sweeps: List[Dict[str, Any]] = field(default_factory=list) # [{"points", "time_lag_s", "age_index"}]
dropped_sweeps: int = 0 # cached sweeps dropped by the capacity rule
def to_dict(self) -> Dict[str, Any]:
return {"num_points": self.num_points, "sweeps": list(self.sweeps), "dropped_sweeps": self.dropped_sweeps}
def _sweep_rows(xyz: np.ndarray, time_lag: np.float32, extra: Optional[np.ndarray], point_feature_size: int,
affine: Optional[np.ndarray]) -> np.ndarray:
"""Rows of one sweep in the network layout: (x, y, z, time_lag) or (x, y, z, intensity, time_lag)."""
out = np.empty((len(xyz), point_feature_size), np.float32)
out[:, :3] = xyz if affine is None else transform_points_f32(xyz, affine)
if point_feature_size == 4:
out[:, 3] = time_lag
elif point_feature_size == 5:
out[:, 3] = 0.0 if extra is None else extra
out[:, 4] = time_lag
else:
raise ValueError(f"point_feature_size must be 4 or 5, got {point_feature_size}")
return out
@dataclass
class _Cached:
xyz: np.ndarray # (N, 3) float32 in the frame of its own stamp
extra: Optional[np.ndarray] # (N,) float32 intensity (point_feature_size 5) or None
stamp_s: float
past2world: np.ndarray # (4, 4) float32 = Affine3f inverse of world2current at its stamp
class SweepDensifier:
"""One stream's Autoware densification cache (module docstring). ``num_past_frames`` is
``densification_params.num_past_frames`` (1 in Autoware's CenterPoint / TransFusion configs)."""
def __init__(self, num_past_frames: int = 1, cloud_capacity: int = DEFAULT_CLOUD_CAPACITY):
if num_past_frames < 0:
raise ValueError("num_past_frames must be >= 0")
self.num_past_frames = int(num_past_frames)
self.cloud_capacity = int(cloud_capacity)
self._cache: List[_Cached] = []
self._world2current = np.eye(4, dtype=np.float32)
self._stamp: Optional[float] = None
@property
def cache_size(self) -> int:
"""``pointcloud_cache_size()`` = past frames + the current one."""
return self.num_past_frames + 1
@property
def needs_pose(self) -> bool:
"""True when enqueue needs ``T_world_from_sensor`` (Autoware looks up a TF only for cache size > 1)."""
return self.cache_size > 1
def __len__(self) -> int:
return len(self._cache)
@property
def stamps(self) -> List[float]:
"""Stamps of the cached sweeps, newest first."""
return [c.stamp_s for c in self._cache]
def reset(self) -> None:
self._cache.clear()
self._world2current = np.eye(4, dtype=np.float32)
self._stamp = None
def enqueue(self, xyz: np.ndarray, stamp_s: float, T_world_from_sensor: Optional[np.ndarray] = None, *,
intensity: Optional[np.ndarray] = None) -> bool:
"""``PointCloudDensification::enqueuePointCloud``: returns False, leaving the cache unchanged, when a pose is
needed and missing (Autoware's TF failure: the frame is skipped). ``xyz`` (N, >=3) is in the sensor frame
of ``stamp_s`` (base_link for Autoware's concatenated cloud); ``T_world_from_sensor`` is that frame's pose
in the world (``map``) frame, float64."""
pts = np.ascontiguousarray(np.asarray(xyz, dtype=np.float32)[:, :3])
extra = None if intensity is None else np.asarray(intensity, dtype=np.float32).reshape(-1)
if extra is not None and len(extra) != len(pts):
raise ValueError("intensity length differs from the point count")
if self.needs_pose:
if T_world_from_sensor is None:
return False
# tf2 lookupTransform(frame <- world) in double, then .cast<float>() (pointcloud_densification.cpp:47-51)
world2current = invert_rigid(np.asarray(T_world_from_sensor, dtype=np.float64)).astype(np.float32)
else:
world2current = np.eye(4, dtype=np.float32)
self._world2current = world2current
self._stamp = float(stamp_s)
self._cache.insert(0, _Cached(pts, extra, float(stamp_s), invert_affine_f32(world2current)))
if len(self._cache) > self.cache_size:
self._cache.pop()
return True
def sweep_points(self, *, point_feature_size: int = 4) -> Tuple[np.ndarray, DensifyInfo]:
"""``VoxelGenerator::generateSweepPoints``: every cached sweep in the current frame, newest first, as
(N, point_feature_size) float32 rows (x, y, z[, intensity], time_lag)."""
info = DensifyInfo()
if not self._cache:
return np.zeros((0, point_feature_size), np.float32), info
parts: List[np.ndarray] = []
n = 0
for k, c in enumerate(self._cache):
if n + len(c.xyz) > self.cloud_capacity:
info.dropped_sweeps = len(self._cache) - k
break
affine = compose_f32(self._world2current, c.past2world)
time_lag = np.float32(self._stamp - c.stamp_s) # double difference cast to float
parts.append(_sweep_rows(c.xyz, time_lag, c.extra, point_feature_size, affine))
info.sweeps.append({"points": int(len(c.xyz)), "time_lag_s": float(time_lag), "age_index": k})
n += len(c.xyz)
info.num_points = n
out = np.concatenate(parts, axis=0) if parts else np.zeros((0, point_feature_size), np.float32)
return out, info
def describe(self) -> Dict[str, Any]:
return {"num_past_frames": self.num_past_frames, "cloud_capacity": self.cloud_capacity,
"cached": len(self._cache), "stamps": self.stamps}
class StreamDensifiers:
"""Per-stream densification caches (PLAN.md D16: host sweep caches are a per-``stream.id`` dict). A new id beyond
``max_streams`` evicts the least recently used stream; ``get(id, reset=True)`` restarts one (cold start)."""
def __init__(self, num_past_frames: int = 1, cloud_capacity: int = DEFAULT_CLOUD_CAPACITY,
max_streams: int = 16):
if max_streams < 1:
raise ValueError("max_streams must be >= 1")
self.num_past_frames = int(num_past_frames)
self.cloud_capacity = int(cloud_capacity)
self.max_streams = int(max_streams)
self._streams: "collections.OrderedDict[str, SweepDensifier]" = collections.OrderedDict()
def get(self, stream_id: str = "default", *, reset: bool = False) -> SweepDensifier:
sid = str(stream_id)
d = self._streams.get(sid)
if d is None:
d = SweepDensifier(self.num_past_frames, self.cloud_capacity)
self._streams[sid] = d
while len(self._streams) > self.max_streams:
self._streams.popitem(last=False)
elif reset:
d.reset()
self._streams.move_to_end(sid)
return d
def forget(self, stream_id: str) -> None:
self._streams.pop(str(stream_id), None)
@property
def streams(self) -> List[str]:
"""Stream ids, most recently used first."""
return list(reversed(self._streams))
def describe(self) -> Dict[str, Any]:
return {"max_streams": self.max_streams, "num_past_frames": self.num_past_frames,
"streams": {k: v.describe() for k, v in reversed(self._streams.items())}}
def densify_sweeps(current_xyz: np.ndarray, sweeps: Sequence[Tuple[np.ndarray, float, Optional[np.ndarray]]] = (),
*, point_feature_size: int = 4, cloud_capacity: int = DEFAULT_CLOUD_CAPACITY,
current_intensity: Optional[np.ndarray] = None,
sweep_intensity: Optional[Sequence[Optional[np.ndarray]]] = None) -> Tuple[np.ndarray, DensifyInfo]:
"""Stateless densification for clients that hold the past sweeps: the current sweep (identity, ``time_lag``
0) then each ``(xyz, time_lag_s, T_current_from_sweep)`` in the given order (newest first, as Autoware's cache),
transformed by the float32 sweep kernel (``T`` None = already in the current frame) with
``time_lag = float32(time_lag_s)``; the capacity rule drops a sweep and every later one."""
info = DensifyInfo()
cur = np.asarray(current_xyz, dtype=np.float32)[:, :3]
rows = [(cur, np.float32(0.0), current_intensity, None)]
for i, (xyz, lag, T) in enumerate(sweeps):
extra = None if sweep_intensity is None else sweep_intensity[i]
rows.append((np.asarray(xyz, dtype=np.float32)[:, :3], np.float32(lag), extra,
None if T is None else np.asarray(T, dtype=np.float64)))
parts = []
n = 0
for k, (xyz, lag, extra, T) in enumerate(rows):
if n + len(xyz) > cloud_capacity:
info.dropped_sweeps = len(rows) - k
break
parts.append(_sweep_rows(xyz, lag, extra, point_feature_size, None if T is None else T.astype(np.float32)))
info.sweeps.append({"points": int(len(xyz)), "time_lag_s": float(lag), "age_index": k})
n += len(xyz)
info.num_points = n
out = np.concatenate(parts, axis=0) if parts else np.zeros((0, point_feature_size), np.float32)
return out, info
# -------------------------------------------------------------------------------------------------- hygiene
def finite_rows(points: np.ndarray, columns: Sequence[int] = (0, 1, 2)) -> np.ndarray:
"""Mask of rows whose ``columns`` are all finite."""
p = np.asarray(points)
return np.isfinite(p[:, list(columns)]).all(axis=1)
def nonzero_rows(points: np.ndarray) -> np.ndarray:
"""Mask of rows whose x, y, z are not all exactly zero (organized clouds mark missing returns as (0, 0, 0); the
research references drop them with ``|x| + |y| + |z| > 0``, which also drops non-finite rows)."""
p = np.asarray(points, dtype=np.float32)
return np.abs(p[:, :3]).sum(axis=1) > 0
def range_mask(points: np.ndarray, range_min: Sequence[float], range_max: Sequence[float]) -> np.ndarray:
"""Autoware's half-open box test ``min <= v < max`` per axis on x, y, z (float32 comparisons). Unlike Autoware's
rejection form, a NaN coordinate fails it."""
p = np.asarray(points, dtype=np.float32)
lo = np.asarray(range_min, dtype=np.float32)
hi = np.asarray(range_max, dtype=np.float32)
keep = np.ones(len(p), dtype=bool)
for a in range(3):
keep &= (p[:, a] >= lo[a]) & (p[:, a] < hi[a])
return keep
def seeded_permutation(n: int, seed: Optional[int]) -> Optional[np.ndarray]:
"""The deterministic stand-in for Autoware's random point shuffle: ``default_rng(seed).permutation(n)`` (PCG64,
numpy >= 1.17), or None when ``seed`` is None (input order)."""
if seed is None:
return None
return np.random.default_rng(int(seed)).permutation(int(n))