Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions src/hflow/checks.py
Original file line number Diff line number Diff line change
Expand Up @@ -205,8 +205,9 @@ def _timestamp_regularity_keys(
(``median_dt_s``/``period_violation_pct``/``max_gap_s``). Across all
selected topics: a pair of ``sync/<cam>~<ref>/{start,end}_offset_s``
keys per populated camera when the episode carries both camera and state
streams -- the densest non-camera stream is the reference. Empty camera
channels retain their per-topic sample count but have no sync offsets. ``App``'s
streams -- the densest populated non-camera stream is the reference. Empty camera
channels retain their per-topic sample count but have no sync offsets. Candidate state
streams with no timestamps are similarly ignored for sync reference. ``App``'s
pre-decode supersession consults this function through the routing
map, which only ever sees the automatic bare registration. The
selection rule mirrors the body's exactly: ``topics=`` is taken as
Expand All @@ -226,6 +227,7 @@ def _timestamp_regularity_keys(
for topic in episode.cameras
if topic in selected and episode.channel(topic).timestamps.size > 0
]
state_topics = [topic for topic in state_topics if episode.channel(topic).timestamps.size > 0]
if camera_topics and state_topics:
reference = max(state_topics, key=lambda topic: episode.topics[topic].message_count)
for camera in camera_topics:
Expand Down Expand Up @@ -320,7 +322,7 @@ async def timestamp_regularity(
same way post-hoc with a far tighter default tolerance; raw multi-sensor
capture needs the looser default here). Deltas beyond ``gap_factor``
periods become labeled gap intervals. Cross-stream: start/end offsets of
every populated camera stream against the densest non-camera stream.
every populated camera stream against the densest populated non-camera stream.

The emitted key set is owned by :func:`_timestamp_regularity_keys`: this
body iterates that function's output and routes each key through
Expand Down Expand Up @@ -375,6 +377,7 @@ def _measure_timestamp_regularity(
for topic in episode.cameras
if topic in selected and per_topic[topic].stamps_ns.size > 0
]
state_topics = [topic for topic in state_topics if per_topic[topic].stamps_ns.size > 0]
if camera_topics and state_topics:
reference = max(state_topics, key=lambda topic: infos[topic].message_count)
sync = _TimestampRegularitySync(
Expand Down
71 changes: 71 additions & 0 deletions tests/test_checks.py
Original file line number Diff line number Diff line change
Expand Up @@ -257,6 +257,77 @@ def write_statistics(statistics: Statistics, builder: RecordBuilder) -> None:
assert not any("empty_cam" in key for key in result.measurements if key.startswith("sync/"))


def test_timestamp_regularity_ignores_empty_reference_for_sync_offsets(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
from foxglove_schemas_protobuf.CompressedVideo_pb2 import CompressedVideo
from mcap_protobuf.schema import build_file_descriptor_set

path = tmp_path / "empty_reference.mcap"
with path.open("wb") as stream:
writer = StockWriter(stream, chunk_size=64 * 1024, compression=CompressionType.NONE)
writer.start()
schema_id = writer.register_schema(
name="foxglove.CompressedVideo",
encoding="protobuf",
data=build_file_descriptor_set(CompressedVideo).SerializeToString(),
)
camera = writer.register_channel(
topic="/cam1", schema_id=schema_id, message_encoding="protobuf"
)
empty_joints = writer.register_channel(
topic="/joint_states", schema_id=0, message_encoding="json"
)
populated_wrench = writer.register_channel(
topic="/wrench", schema_id=0, message_encoding="json"
)

# MCAP summary statistics report message_count >= 2 for the empty state
# stream, but 0 actual records were written to disk.
original_write = Statistics.write

def write_statistics(statistics: Statistics, builder: RecordBuilder) -> None:
original_write(
replace(
statistics,
channel_message_counts={camera: 2, empty_joints: 10, populated_wrench: 2},
),
builder,
)

monkeypatch.setattr(Statistics, "write", write_statistics)
for timestamp_ns in (1_000_000_000, 2_000_000_000):
message = CompressedVideo()
message.timestamp.FromNanoseconds(timestamp_ns)
message.data = b"x"
message.format = "h264"
writer.add_message(camera, timestamp_ns, message.SerializeToString(), timestamp_ns)
writer.add_message(populated_wrench, timestamp_ns, b'{"fx": 0}', timestamp_ns)
writer.finish()

# Case 1: An empty reference state stream is skipped in favor of a populated one
with hflow.Episode(path) as episode:
assert episode.topics["/joint_states"].message_count == 10
assert episode.channel("/joint_states").timestamps.size == 0
result = asyncio.run(
timestamp_regularity(episode, topics=["/cam1", "/joint_states", "/wrench"])
)

assert result.measurements["/joint_states/period_sample_count"] == 0
assert result.measurements["sync//cam1~/wrench/start_offset_s"] == 0.0
assert result.measurements["sync//cam1~/wrench/end_offset_s"] == 0.0
assert not any("joint_states" in key for key in result.measurements if key.startswith("sync/"))

# Case 2: When all state streams are empty, no sync offset keys are produced
with hflow.Episode(path) as episode:
result_no_state = asyncio.run(
timestamp_regularity(episode, topics=["/cam1", "/joint_states"])
)

assert result_no_state.measurements["/joint_states/period_sample_count"] == 0
assert not any(key.startswith("sync/") for key in result_no_state.measurements)


def test_joint_discontinuity_finds_the_injected_jump(jittery_episode: hflow.Episode) -> None:
result = asyncio.run(joint_discontinuity(jittery_episode, velocity_limit=3.0))
violation_count = result.measurements["/joint_states/violation_count"]
Expand Down
Loading