Move analysis and saving out of async
This commit is contained in:
parent
817c53d1df
commit
1ea308a42d
1 changed files with 78 additions and 90 deletions
|
|
@ -268,67 +268,86 @@ class BaseCamera(OFMThing):
|
||||||
self.lores_mjpeg_stream._streaming = False
|
self.lores_mjpeg_stream._streaming = False
|
||||||
|
|
||||||
async def _monitor_framerate(
|
async def _monitor_framerate(
|
||||||
self, duration: float, output_dir: str, sample_interval: float = 0.1
|
self,
|
||||||
) -> None:
|
duration: float,
|
||||||
|
sample_interval: float = 0.1,
|
||||||
|
) -> tuple[float, int, list[dict[str, float | int]]]:
|
||||||
|
"""Asynchronously monitor the timing on incoming frames."""
|
||||||
|
start_time = time.time()
|
||||||
|
last_sample_time = start_time
|
||||||
|
last_sample_frames = 0
|
||||||
|
|
||||||
|
frames = 0
|
||||||
|
samples = []
|
||||||
|
|
||||||
|
async for frame in self.mjpeg_stream.frame_async_generator():
|
||||||
|
if not self._framerate_monitor_running:
|
||||||
|
break
|
||||||
|
|
||||||
|
now = time.time()
|
||||||
|
frames += 1
|
||||||
|
|
||||||
|
if now - last_sample_time >= sample_interval:
|
||||||
|
interval = now - last_sample_time
|
||||||
|
fps = (frames - last_sample_frames) / interval
|
||||||
|
|
||||||
|
samples.append(
|
||||||
|
{
|
||||||
|
"timestamp": now,
|
||||||
|
"frame_count": frames,
|
||||||
|
"frame_size_bytes": len(frame),
|
||||||
|
"instant_fps": fps,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
last_sample_time = now
|
||||||
|
last_sample_frames = frames
|
||||||
|
|
||||||
|
if now - start_time >= duration:
|
||||||
|
break
|
||||||
|
|
||||||
|
total_time = time.time() - start_time
|
||||||
|
|
||||||
|
return total_time, frames, samples
|
||||||
|
|
||||||
|
@lt.action
|
||||||
|
def record_framerate(
|
||||||
|
self,
|
||||||
|
duration: float = 5.0,
|
||||||
|
) -> str:
|
||||||
|
"""Record MJPEG stream framerate statistics."""
|
||||||
|
output_dir = os.path.join(
|
||||||
|
self.data_dir,
|
||||||
|
"characterisation",
|
||||||
|
"framerate",
|
||||||
|
)
|
||||||
|
os.makedirs(output_dir, exist_ok=True)
|
||||||
|
|
||||||
|
timestamp = time.strftime("%Y%m%d_%H%M%S")
|
||||||
|
log_path = os.path.join(
|
||||||
|
output_dir,
|
||||||
|
f"framerate_{timestamp}.json",
|
||||||
|
)
|
||||||
|
|
||||||
|
self.logger.info(
|
||||||
|
"Framerate monitor started -> %s",
|
||||||
|
log_path,
|
||||||
|
)
|
||||||
|
|
||||||
|
self._framerate_monitor_running = True
|
||||||
|
|
||||||
|
# This runs as an async task, which we wait to complete
|
||||||
try:
|
try:
|
||||||
os.makedirs(output_dir, exist_ok=True)
|
total_time, frames, samples = (
|
||||||
|
self._thing_server_interface.call_async_task(
|
||||||
timestamp = time.strftime("%Y%m%d_%H%M%S")
|
self._monitor_framerate,
|
||||||
log_path = os.path.join(output_dir, f"framerate_{timestamp}.json")
|
duration,
|
||||||
|
)
|
||||||
start_time = time.time()
|
|
||||||
last_sample_time = start_time
|
|
||||||
last_sample_frames = 0
|
|
||||||
|
|
||||||
frames = 0
|
|
||||||
samples = []
|
|
||||||
|
|
||||||
self._framerate_monitor_running = True
|
|
||||||
|
|
||||||
self.logger.info("Framerate monitor started -> %s", log_path)
|
|
||||||
|
|
||||||
async for frame in self.mjpeg_stream.frame_async_generator():
|
|
||||||
if not self._framerate_monitor_running:
|
|
||||||
break
|
|
||||||
|
|
||||||
now = time.time()
|
|
||||||
frames += 1
|
|
||||||
|
|
||||||
# take data every sample_interval seconds
|
|
||||||
if now - last_sample_time >= sample_interval:
|
|
||||||
interval = now - last_sample_time
|
|
||||||
fps = (frames - last_sample_frames) / interval
|
|
||||||
|
|
||||||
samples.append((now, frames, len(frame), fps))
|
|
||||||
|
|
||||||
last_sample_time = now
|
|
||||||
last_sample_frames = frames
|
|
||||||
|
|
||||||
if now - start_time >= duration:
|
|
||||||
break
|
|
||||||
|
|
||||||
# write to disk outside of async
|
|
||||||
await asyncio.to_thread(
|
|
||||||
self._write_framerate_json, log_path, samples, start_time, frames
|
|
||||||
)
|
)
|
||||||
|
|
||||||
self.logger.info("Framerate monitor complete")
|
|
||||||
|
|
||||||
except Exception:
|
|
||||||
self.logger.exception("Framerate monitor crashed")
|
|
||||||
|
|
||||||
finally:
|
finally:
|
||||||
self._framerate_monitor_running = False
|
self._framerate_monitor_running = False
|
||||||
|
|
||||||
def _write_framerate_json(
|
|
||||||
self,
|
|
||||||
log_path: str,
|
|
||||||
samples: list[tuple[float, int, int, float]],
|
|
||||||
start_time: float,
|
|
||||||
frames: int,
|
|
||||||
) -> None:
|
|
||||||
"""Write framerate monitoring results to a JSON file."""
|
|
||||||
total_time = time.time() - start_time
|
|
||||||
avg_fps = frames / total_time if total_time > 0 else 0
|
avg_fps = frames / total_time if total_time > 0 else 0
|
||||||
|
|
||||||
data = {
|
data = {
|
||||||
|
|
@ -337,49 +356,18 @@ class BaseCamera(OFMThing):
|
||||||
"total_frames": frames,
|
"total_frames": frames,
|
||||||
"avg_fps": avg_fps,
|
"avg_fps": avg_fps,
|
||||||
},
|
},
|
||||||
"samples": [
|
"samples": samples,
|
||||||
{
|
|
||||||
"timestamp": s[0],
|
|
||||||
"frame_count": s[1],
|
|
||||||
"frame_size_bytes": s[2],
|
|
||||||
"instant_fps": s[3],
|
|
||||||
}
|
|
||||||
for s in samples
|
|
||||||
],
|
|
||||||
}
|
}
|
||||||
|
|
||||||
with open(log_path, "w") as f:
|
with open(log_path, "w") as f:
|
||||||
json.dump(data, f, indent=2)
|
json.dump(data, f, indent=2)
|
||||||
|
|
||||||
@lt.action
|
|
||||||
def record_framerate_async(
|
|
||||||
self,
|
|
||||||
duration: float = 5.0,
|
|
||||||
) -> str:
|
|
||||||
"""Start asynchronous MJPEG framerate monitoring.
|
|
||||||
|
|
||||||
The monitor observes frames delivered through ``self.mjpeg_stream`` rather
|
|
||||||
than directly querying the camera hardware. This makes it suitable for
|
|
||||||
measuring effective stream throughput under real operating conditions.
|
|
||||||
|
|
||||||
:param duration: Monitoring duration in seconds. Default = 5.0s.
|
|
||||||
|
|
||||||
:returns: The output directory containing generated log files.
|
|
||||||
"""
|
|
||||||
output_dir = os.path.join(self.data_dir, "characterisation", "framerate")
|
|
||||||
self._thing_server_interface.start_async_task_soon(
|
|
||||||
self._monitor_framerate,
|
|
||||||
duration,
|
|
||||||
output_dir,
|
|
||||||
)
|
|
||||||
|
|
||||||
self.logger.info(
|
self.logger.info(
|
||||||
"Started async framerate monitor for %.2fs -> %s",
|
"Framerate monitor complete -> %s",
|
||||||
duration,
|
log_path,
|
||||||
output_dir,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
return output_dir
|
return log_path
|
||||||
|
|
||||||
@lt.property
|
@lt.property
|
||||||
def stream_active(self) -> bool:
|
def stream_active(self) -> bool:
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue