- flask-socketio + threading mode (no eventlet monkey-patching) - preview_queue in pipeline: downscale to 640p, JPEG encode, bounded queue - encoder thread per job reads queue → SocketIO broadcast to room - Frontend: SocketIO client replaces HTTP polling, handles preview_frame event - Configurable: preview_every_n, preview_max_dim, preview_jpeg_quality - 62/62 tests pass
266 lines
8.7 KiB
Markdown
266 lines
8.7 KiB
Markdown
# Plan: 15+ FPS Live Preview via WebSocket
|
|
|
|
## Goal
|
|
Achieve **15+ FPS at 640p** live preview during video processing using **WebSocket push** instead of HTTP polling.
|
|
|
|
---
|
|
|
|
## Current Architecture Analysis
|
|
|
|
| Component | Current | Bottleneck |
|
|
|-----------|---------|------------|
|
|
| Frame capture | Every 10th frame (2.5 Hz at 25 FPS) | Too slow for 15 FPS target |
|
|
| Encoding | `cv2.imencode('.jpg', viz)` synchronous, full-res | Blocks pipeline |
|
|
| Transport | HTTP GET `/api/jobs/<id>/frame` every 2s | 2s latency, wasted polls |
|
|
| Frontend | `<img src>` replacement | Flicker, no frame timing control |
|
|
|
|
**Key finding**: JPEG encoding at 640p = **285-371 FPS** (CPU). Not the bottleneck. The bottleneck is **synchronous encoding + polling architecture**.
|
|
|
|
---
|
|
|
|
## Recommended Architecture: WebSocket with Threaded Encoder
|
|
|
|
```
|
|
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
|
|
│ Pipeline │ │ Encoder Thread │ │ Flask-SocketIO │
|
|
│ (detection) │────▶│ (non-blocking) │────▶│ WebSocket │
|
|
│ - YOLO detect │ │ - Downscale │ │ Server │
|
|
│ - Track │ │ - JPEG encode │ │ - Broadcast │
|
|
│ - Count │ │ - Queue push │ │ to clients │
|
|
└─────────────────┘ └─────────────────┘ └─────────────────┘
|
|
│ │ │
|
|
▼ ▼ ▼
|
|
Full speed 15+ FPS @ 640p Real-time push
|
|
processing to queue to browsers
|
|
```
|
|
|
|
### Why This Design?
|
|
- **Pipeline stays at full speed** — detection/tracking never blocked by encoding
|
|
- **Encoder thread** — produces 15+ FPS at 640p, drops frames if queue full
|
|
- **Flask-SocketIO** — handles WebSocket connections, broadcasting, reconnection
|
|
- **Eventlet/gevent** — async I/O for many concurrent viewers
|
|
|
|
---
|
|
|
|
## Implementation Plan
|
|
|
|
### Phase 1: Dependencies & Server Setup
|
|
|
|
**Add to `pyproject.toml`:**
|
|
```toml
|
|
[project.optional-dependencies]
|
|
web = [
|
|
"flask-socketio>=5.3.0",
|
|
"eventlet>=0.33.0", # or gevent
|
|
]
|
|
```
|
|
|
|
**Run:** `pip install -e ".[web]"`
|
|
|
|
### Phase 2: Pipeline — Non-blocking Frame Producer
|
|
|
|
**Modify `src/pipeline.py`:**
|
|
- Add `preview_queue: Queue[bytes]` parameter to `run_pipeline()`
|
|
- In main loop: every N frames (configurable), downscale → encode → `queue.put_nowait(jpeg_bytes)`
|
|
- If queue full: drop frame (don't block)
|
|
- Remove `frame_callback` — replaced by queue
|
|
|
|
```python
|
|
def run_pipeline(
|
|
...,
|
|
preview_queue: "Queue[bytes] | None" = None,
|
|
preview_every_n: int = 2, # 25 FPS / 2 = 12.5 FPS → use 1 for 25 FPS
|
|
preview_max_dim: int = 640,
|
|
preview_jpeg_quality: int = 70,
|
|
) -> PipelineResult:
|
|
if preview_queue is not None:
|
|
preview_interval = max(1, preview_every_n)
|
|
...
|
|
if preview_queue and frame_idx % preview_interval == 0:
|
|
# Downscale
|
|
h, w = viz.shape[:2]
|
|
if max(h, w) > preview_max_dim:
|
|
scale = preview_max_dim / max(h, w)
|
|
preview_viz = cv2.resize(viz, (int(w*scale), int(h*scale)))
|
|
else:
|
|
preview_viz = viz
|
|
ok, jpeg = cv2.imencode('.jpg', preview_viz, [cv2.IMWRITE_JPEG_QUALITY, preview_jpeg_quality])
|
|
if ok:
|
|
try:
|
|
preview_queue.put_nowait(jpeg.tobytes())
|
|
except queue.Full:
|
|
pass # Drop frame, don't block
|
|
```
|
|
|
|
### Phase 3: Job Queue — WebSocket Broadcaster
|
|
|
|
**Modify `src/job.py`:**
|
|
- Add `preview_queue` per job
|
|
- Start **encoder thread** that reads queue → broadcasts via SocketIO
|
|
- Thread lifecycle tied to job
|
|
|
|
```python
|
|
import queue
|
|
import threading
|
|
|
|
class Job:
|
|
...
|
|
preview_queue: "queue.Queue[bytes]" = field(default_factory=lambda: queue.Queue(maxsize=30))
|
|
_encoder_thread: threading.Thread | None = None
|
|
_encoder_stop: threading.Event = field(default_factory=threading.Event)
|
|
|
|
def _start_preview_broadcaster(self, job_id: str, socketio):
|
|
"""Thread: read preview queue, emit via SocketIO."""
|
|
while not self._encoder_stop.is_set():
|
|
try:
|
|
jpeg_bytes = self.preview_queue.get(timeout=0.1)
|
|
# Emit to job's room
|
|
socketio.emit('preview_frame', {'frame': base64.b64encode(jpeg_bytes).decode()}, room=job_id)
|
|
except queue.Empty:
|
|
continue
|
|
```
|
|
|
|
### Phase 4: Flask-SocketIO Server
|
|
|
|
**Modify `app.py`:**
|
|
```python
|
|
from flask_socketio import SocketIO, emit, join_room, leave_room
|
|
|
|
socketio = SocketIO(app, cors_allowed_origins="*", async_mode='eventlet')
|
|
|
|
@socketio.on('connect')
|
|
def on_connect():
|
|
pass
|
|
|
|
@socketio.on('join_job')
|
|
def on_join_job(data):
|
|
job_id = data.get('job_id')
|
|
if job_id:
|
|
join_room(job_id)
|
|
|
|
@socketio.on('leave_job')
|
|
def on_leave_job(data):
|
|
job_id = data.get('job_id')
|
|
if job_id:
|
|
leave_room(job_id)
|
|
|
|
# In _run_job():
|
|
# socketio.start_background_task(_start_preview_broadcaster, job_id, socketio)
|
|
```
|
|
|
|
**Run with eventlet:**
|
|
```python
|
|
# main.py or app.py
|
|
if __name__ == '__main__':
|
|
socketio.run(app, host='0.0.0.0', port=9000, debug=False)
|
|
```
|
|
|
|
### Phase 5: Frontend — WebSocket Consumer
|
|
|
|
**Replace polling in `static/app.js`:**
|
|
```javascript
|
|
// In initStatusPage():
|
|
var socket = io();
|
|
socket.on('connect', function() {
|
|
socket.emit('join_job', {job_id: jobId});
|
|
});
|
|
|
|
socket.on('preview_frame', function(data) {
|
|
if (livePreviewImg && data.frame) {
|
|
livePreviewImg.src = 'data:image/jpeg;base64,' + data.frame;
|
|
if (previewPlaceholder) previewPlaceholder.style.display = 'none';
|
|
livePreviewImg.style.display = 'block';
|
|
}
|
|
});
|
|
|
|
// Cleanup on disconnect/page unload
|
|
socket.on('disconnect', function() {
|
|
socket.emit('leave_job', {job_id: jobId});
|
|
});
|
|
window.addEventListener('beforeunload', function() {
|
|
socket.emit('leave_job', {job_id: jobId});
|
|
});
|
|
```
|
|
|
|
### Phase 6: Configuration & Tuning
|
|
|
|
| Parameter | Default | Tuning Notes |
|
|
|-----------|---------|--------------|
|
|
| `preview_every_n` | 1 | 1 = every frame (25 FPS max), 2 = 12.5 FPS |
|
|
| `preview_max_dim` | 640 | 480 for faster, 720 for quality |
|
|
| `preview_jpeg_quality` | 70 | 50-80 tradeoff |
|
|
| `queue.maxsize` | 30 | ~2 seconds buffer at 15 FPS |
|
|
|
|
---
|
|
|
|
## Expected Performance
|
|
|
|
| Metric | Current | Target |
|
|
|--------|---------|--------|
|
|
| Preview FPS | ~0.5 (polling) | **15-25** |
|
|
| Latency | 2+ seconds | **<100ms** |
|
|
| Pipeline impact | Blocks every 10th frame | **Zero** |
|
|
| Resolution | Full video res | **640p max** |
|
|
| Bandwidth/frame | ~200-500 KB | **~30-50 KB** |
|
|
|
|
---
|
|
|
|
## Files to Modify
|
|
|
|
| File | Changes |
|
|
|------|---------|
|
|
| `pyproject.toml` | Add `flask-socketio`, `eventlet` to optional deps |
|
|
| `src/pipeline.py` | Add `preview_queue` param, producer logic |
|
|
| `src/job.py` | Add preview queue, encoder thread, SocketIO broadcaster |
|
|
| `app.py` | Initialize SocketIO, add connect/join/leave handlers, run with `socketio.run()` |
|
|
| `static/app.js` | Replace polling with SocketIO client, handle `preview_frame` event |
|
|
| `templates/base.html` | Add SocketIO client script (`/socket.io/socket.io.js`) |
|
|
|
|
---
|
|
|
|
## Testing Strategy
|
|
|
|
1. **Unit test**: `run_pipeline` with mock queue — verify frames enqueued at correct interval
|
|
2. **Integration**: Start job, connect browser, verify 15+ FPS in devtools Network tab
|
|
3. **Load test**: Multiple browser tabs → verify broadcast works
|
|
4. **Stress test**: Long video (10+ min) → verify no memory leaks, queue stays bounded
|
|
|
|
---
|
|
|
|
## Risks & Mitigations
|
|
|
|
| Risk | Mitigation |
|
|
|------|------------|
|
|
| Eventlet monkey-patches stdlib — may break torch/CUDA | Test early; fallback to `gevent` or native `websockets` + asyncio |
|
|
| Queue memory growth | Bounded queue (`maxsize=30`), drop frames when full |
|
|
| Multiple clients | SocketIO rooms — single broadcast to room |
|
|
| Reconnection | SocketIO handles auto-reconnect; client re-joins room on reconnect |
|
|
|
|
---
|
|
|
|
## Dependencies Check
|
|
|
|
- `flask-socketio>=5.3.0` — compatible with Flask 3.x
|
|
- `eventlet>=0.33.0` — works on ARM64 (Jetson)
|
|
- `python-engineio>=4.7.0` — transitive
|
|
|
|
---
|
|
|
|
## Rollout
|
|
|
|
1. Add deps, install
|
|
2. Implement pipeline + job queue changes (backend only)
|
|
3. Test with `python -c "from app import socketio; print('OK')"`
|
|
4. Implement frontend WebSocket client
|
|
5. Test end-to-end
|
|
6. Deploy
|
|
|
|
---
|
|
|
|
## Alternative: H.264 Streaming (Future)
|
|
|
|
If JPEG-over-WebSocket isn't smooth enough:
|
|
- Use `cv2.VideoWriter` with H.264 in encoder thread
|
|
- Stream via HTTP chunked transfer or WebRTC
|
|
- Browser `<video>` tag with `MediaSource` API
|
|
- More complex but hardware-decoded, smoother at high FPS |