feat: switch live preview from WebSocket to MJPEG streaming

- Pipeline writes JPEG to /tmp/feedmill_preview_{job_id}.jpg (atomic)
- Flask serves MJPEG stream at /api/preview/{job_id}
- Frontend uses native <img src> MJPEG — zero JS needed
- Removed flask-socketio, eventlet, socket.io CDN dependencies
- 33% less bandwidth per frame, native browser decode
- Same proven approach as original karung-counting project
This commit is contained in:
jetson committed 2026-09-22 10:09:22 +07:00
1 parent 1e28eb49e2
commit a5566f3b99
8 files changed
+239 -100

No files matched your search

+187
View File
@@ -0,0 +1,187 @@
# Plan: Switch Live Preview from WebSocket+base64 to MJPEG Streaming
## Goal
Replace the WebSocket+base64 live preview with MJPEG streaming (same approach as the original karung-counting project) for faster, smoother preview with less overhead.
## Why MJPEG is Faster
| Factor | WebSocket+base64 | MJPEG |
|--------|-----------------|-------|
| **Frame size** | JPEG + 33% base64 bloat + JSON wrapper | Raw JPEG bytes |
| **Transport** | JSON over WebSocket | Binary HTTP stream |
| **Decode** | JS sets img.src per frame (triggers decode) | Browser native `<img>` MJPEG decode |
| **Code** | Queue + thread + SocketIO emit + JS handler | File write + HTTP generator |
| **Dependencies** | flask-socketio, JS socket.io client | None (stdlib only) |
## Architecture
```
Pipeline (pipeline.py)
→ write JPEG to /tmp/feedmill_preview_{job_id}.jpg (atomic tmp+replace)
Flask endpoint (app.py)
GET /api/preview/<job_id>
→ reads file in tight loop, yields multipart/x-mixed-replace MJPEG stream
Frontend (status.html)
<img src="/api/preview/{job_id}" />
→ browser handles MJPEG natively, no JS needed
```
## Implementation Plan
### Phase 1: Pipeline — Write to file instead of Queue
**Modify `src/pipeline.py`:**
- Add `preview_path: str | None = None` parameter
- Replace queue put with atomic file write:
```python
if preview_path is not None and frame_idx % max(1, preview_every_n) == 0:
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
tmp_path = preview_path + ".tmp.jpg"
cv2.imwrite(tmp_path, preview_viz, [cv2.IMWRITE_JPEG_QUALITY, preview_jpeg_quality])
os.replace(tmp_path, preview_path)
```
- Remove `preview_queue` parameter (deprecated)
### Phase 2: Job Queue — Compute preview path, remove broadcaster
**Modify `src/job.py`:**
- Remove `preview_queue`, `_preview_thread`, `_preview_stop` fields
- Remove `start_preview_broadcaster()` and `stop_preview_broadcaster()` methods
- Add `preview_path: str` field (computed in `__init__`):
```python
preview_path: str = "" # set in add_job to /tmp/feedmill_preview_{job_id}.jpg
```
- In `add_job()`, set `job.preview_path = f"/tmp/feedmill_preview_{job.job_id}.jpg"`
- In `_run_job()`, pass `preview_path=job.preview_path` to `run_pipeline()`
- In cleanup, delete the preview file:
```python
try:
os.remove(job.preview_path)
except OSError:
pass
```
### Phase 3: Flask — Add MJPEG endpoint, remove SocketIO preview
**Modify `app.py`:**
- Add MJPEG endpoint:
```python
import time as _time
@app.route("/api/preview/<job_id>")
def api_preview(job_id):
job = job_queue.get_job(job_id)
if job is None:
return jsonify({"error": "not found"}), 404
preview_path = job.preview_path
if not preview_path:
return jsonify({"error": "no preview path"}), 404
def generate():
consecutive_fails = 0
MAX_FAILS = 30
while True:
try:
with open(preview_path, "rb") as f:
jpeg = f.read()
consecutive_fails = 0
yield (b"--frame\r\n"
b"Content-Type: image/jpeg\r\n\r\n" + jpeg + b"\r\n")
except FileNotFoundError:
consecutive_fails += 1
if consecutive_fails >= MAX_FAILS:
return
_time.sleep(1.0)
continue
except Exception:
consecutive_fails += 1
if consecutive_fails >= MAX_FAILS:
return
_time.sleep(0.5)
continue
_time.sleep(0.05)
return Response(generate(), mimetype="multipart/x-mixed-replace; boundary=frame")
```
- Remove SocketIO handlers for preview (`on_connect`, `on_disconnect`, `on_join_job`, `on_leave_job`)
- Remove `job.start_preview_broadcaster()` calls from upload routes
- Remove `flask_socketio` import and `socketio` init (if no longer needed)
- Keep `socketio.run()` in `main()` OR switch to `app.run()` if SocketIO is fully removed
### Phase 4: Frontend — Replace SocketIO with MJPEG img tag
**Modify `templates/status.html`:**
- Change `<img>` src to MJPEG endpoint:
```html
<img id="live-preview-img" class="live-preview-img"
src="/api/preview/{{ job.job_id }}"
alt="Annotated frame from video processing" />
```
- Remove `preview-placeholder` div (MJPEG auto-shows when frames arrive)
**Modify `static/app.js`:**
- Remove SocketIO code:
```javascript
// DELETE: var socket = io();
// DELETE: socket.on('connect', ...)
// DELETE: socket.on('preview_frame', ...)
// DELETE: window.addEventListener('beforeunload', ...)
```
- Remove `updateLivePreview()` function (no longer called)
### Phase 5: Cleanup — Remove SocketIO if fully unused
**Modify `app.py`:**
- Remove `from flask_socketio import SocketIO, join_room, leave_room`
- Remove `socketio = SocketIO(...)` line
- Change `socketio.run(app, ...)` to `app.run(host, port, debug)`
**Modify `pyproject.toml`:**
- Remove `flask-socketio` and `eventlet` from optional deps
**Modify `templates/base.html`:**
- Remove SocketIO CDN script tag
## Files to Modify
| File | Changes |
|------|---------|
| `src/pipeline.py` | Replace `preview_queue` with `preview_path` file write |
| `src/job.py` | Compute `preview_path`, remove broadcaster thread/queue |
| `app.py` | Add MJPEG endpoint, remove SocketIO handlers, use `app.run()` |
| `templates/status.html` | Add `src="/api/preview/{job_id}"` to img tag |
| `templates/base.html` | Remove SocketIO CDN script |
| `static/app.js` | Remove SocketIO code and `updateLivePreview()` |
| `pyproject.toml` | Remove `flask-socketio`, `eventlet` deps |
## Expected Performance
| Metric | Before (WebSocket) | After (MJPEG) |
|--------|-------------------|---------------|
| Frame size | ~50-80 KB (base64+JSON) | ~30-50 KB (raw JPEG) |
| Transport overhead | 33% base64 + JSON | 0% |
| Decode latency | JS img.src assignment | Native browser MJPEG |
| Dependencies | flask-socketio, eventlet, socket.io CDN | None |
| JS complexity | ~20 lines SocketIO code | 0 lines |
## Testing
1. Run `python -m pytest tests/ -v --tb=short` — 62/62 pass
2. Start web UI, upload video, start job
3. Verify MJPEG stream loads at `http://localhost:9000/api/preview/{job_id}`
4. Verify live preview updates smoothly in browser
5. Verify cleanup: preview file deleted after job completes
## Rollout
1. Implement all changes
2. Test locally
3. Commit and push
4. Restart web UI
+34 -40
View File
@@ -10,7 +10,8 @@ from flask import (
Flask, render_template, request, redirect,
url_for, send_file, jsonify,
)
from flask_socketio import SocketIO, join_room, leave_room
import time as _time
from flask import Response
from werkzeug.utils import secure_filename
from pathlib import Path
@@ -25,8 +26,6 @@ app = Flask(__name__, template_folder="templates", static_folder="static")
app.config["SECRET_KEY"] = os.getenv("SECRET_KEY", "change-me")
app.config["MAX_CONTENT_LENGTH"] = 2 * 1024 * 1024 * 1024 # 2GB
socketio = SocketIO(app, cors_allowed_origins="*", async_mode='threading')
MODELS_DIR = os.getenv("MODELS_DIR", "./models")
UPLOAD_DIR = os.getenv("UPLOAD_DIR", "./uploads")
OUTPUT_DIR = os.getenv("OUTPUT_DIR", "./output")
@@ -39,30 +38,6 @@ job_queue = JobQueue(output_dir=OUTPUT_DIR)
VIDEO_EXTENSIONS = ('.mp4', '.avi', '.mkv', '.mov', '.webm')
@socketio.on('connect')
def on_connect():
pass
@socketio.on('disconnect')
def on_disconnect():
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)
def _parse_model_configs(form):
"""Parse model selection from request form.
@@ -155,8 +130,6 @@ def upload():
model_configs=model_configs,
class_filters=class_filters,
)
job.start_preview_broadcaster(job.job_id, socketio)
return redirect(url_for("status", job_id=job.job_id))
@@ -176,8 +149,6 @@ def upload_reuse():
model_configs=model_configs,
class_filters=class_filters,
)
job.start_preview_broadcaster(job.job_id, socketio)
return redirect(url_for("status", job_id=job.job_id))
@@ -348,16 +319,39 @@ def api_job_samples(job_id):
return jsonify(rel_paths)
@app.route("/api/jobs/<job_id>/frame")
def api_job_frame(job_id):
"""Return current latest_frame as JPEG (live preview during processing)."""
@app.route("/api/preview/<job_id>")
def api_preview(job_id):
job = job_queue.get_job(job_id)
if job is None:
return "Job not found", 404
if job.latest_frame is None:
return "No frame available", 404
from flask import Response
return Response(job.latest_frame, mimetype="image/jpeg")
return jsonify({"error": "not found"}), 404
preview_path = job.preview_path
if not preview_path:
return jsonify({"error": "no preview path"}), 404
def generate():
consecutive_fails = 0
MAX_FAILS = 30
while True:
try:
with open(preview_path, "rb") as f:
jpeg = f.read()
consecutive_fails = 0
yield (b"--frame\r\n"
b"Content-Type: image/jpeg\r\n\r\n" + jpeg + b"\r\n")
except FileNotFoundError:
consecutive_fails += 1
if consecutive_fails >= MAX_FAILS:
return
_time.sleep(1.0)
continue
except Exception:
consecutive_fails += 1
if consecutive_fails >= MAX_FAILS:
return
_time.sleep(0.5)
continue
_time.sleep(0.05)
return Response(generate(), mimetype="multipart/x-mixed-replace; boundary=frame")
def main():
@@ -366,7 +360,7 @@ def main():
debug = os.getenv("FLASK_DEBUG", "false").lower() == "true"
print(f"Feedmill Recounter web UI: http://{host}:{port}")
socketio.run(app, host=host, port=port, debug=debug, allow_unsafe_werkzeug=True)
app.run(host=host, port=port, debug=debug)
if __name__ == "__main__":
-1
View File
@@ -18,7 +18,6 @@ dependencies = [
[project.optional-dependencies]
dev = ["pytest"]
web = ["flask-socketio>=5.3.0", "eventlet>=0.33.0"]
[project.scripts]
recounter = "cli:main"
+9 -28
View File
@@ -3,9 +3,7 @@
from __future__ import annotations
import base64
import os
import queue
import threading
import time
import uuid
@@ -57,31 +55,7 @@ class Job:
created_at: float = field(default_factory=time.time)
completed_at: float | None = None
latest_frame: bytes | None = None
preview_queue: queue.Queue = field(default_factory=lambda: queue.Queue(maxsize=30))
_preview_thread: threading.Thread | None = None
_preview_stop: threading.Event = field(default_factory=threading.Event)
_job_lock: threading.Lock = field(default_factory=threading.Lock)
def start_preview_broadcaster(self, job_id: str, socketio) -> None:
"""Start background thread that reads preview_queue and emits via SocketIO."""
def _broadcaster():
while not self._preview_stop.is_set():
try:
jpeg_bytes = self.preview_queue.get(timeout=0.2)
with self._job_lock:
self.latest_frame = jpeg_bytes
encoded = base64.b64encode(jpeg_bytes).decode('ascii')
socketio.emit('preview_frame', {'frame': encoded}, room=job_id)
except queue.Empty:
continue
except Exception:
continue
self._preview_thread = threading.Thread(target=_broadcaster, daemon=True)
self._preview_thread.start()
def stop_preview_broadcaster(self) -> None:
self._preview_stop.set()
preview_path: str = ""
class JobQueue:
@@ -111,6 +85,8 @@ class JobQueue:
)
Path(job.output_dir).mkdir(parents=True, exist_ok=True)
job.preview_path = f"/tmp/feedmill_preview_{job.job_id}.jpg"
with self._lock:
self._jobs[job_id] = job
@@ -195,7 +171,7 @@ class JobQueue:
output_path=output_path,
class_filter=class_filter,
truck_model_config=truck_model_config,
preview_queue=job.preview_queue,
preview_path=job.preview_path,
preview_every_n=2,
cancel_check=lambda: job.status == JobStatus.CANCELLED,
)
@@ -237,6 +213,11 @@ class JobQueue:
job.latest_frame = None
finally:
try:
os.remove(job.preview_path)
except OSError:
pass
with self._lock:
job.completed_at = time.time()
job.current_model = ""
+8 -11
View File
@@ -2,7 +2,7 @@
from __future__ import annotations
import queue
import os
import time
from collections.abc import Callable
from dataclasses import dataclass
@@ -65,7 +65,7 @@ def run_pipeline(
progress_callback: Callable[[int, int], None] | None = None,
frame_callback: Callable[[bytes], None] | None = None,
cancel_check: Callable[[], bool] | None = None,
preview_queue: queue.Queue[bytes] | None = None,
preview_path: str | None = None,
preview_every_n: int = 2,
preview_max_dim: int = 640,
preview_jpeg_quality: int = 70,
@@ -83,8 +83,8 @@ def run_pipeline(
progress_callback: Optional fn(frame_idx, total_frames) called per frame.
frame_callback: Optional fn(jpeg_bytes) called every 10th frame (fallback).
cancel_check: Optional fn() returning True to abort processing.
preview_queue: Optional queue for live preview frames (JPEG bytes).
preview_every_n: Push a preview frame every N frames (default 2).
preview_path: Optional file path for live preview frames (JPEG written atomically).
preview_every_n: Write a preview frame every N frames (default 2).
preview_max_dim: Max dimension (w or h) for preview frames.
preview_jpeg_quality: JPEG quality for preview frames (1-100).
@@ -242,19 +242,16 @@ def run_pipeline(
writer.write_frame(viz)
if preview_queue is not None and frame_idx % max(1, preview_every_n) == 0:
if preview_path is not None and frame_idx % max(1, preview_every_n) == 0:
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 Exception:
pass
tmp_path = preview_path + ".tmp.jpg"
cv2.imwrite(tmp_path, preview_viz, [cv2.IMWRITE_JPEG_QUALITY, preview_jpeg_quality])
os.replace(tmp_path, preview_path)
elif frame_callback and frame_idx % 10 == 0:
ok, jpeg = cv2.imencode('.jpg', viz)
if ok:
-18
View File
@@ -358,24 +358,6 @@ function initStatusPage(jobId, initialStatus) {
if (!progressFill) return;
/* --- WebSocket live preview --- */
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';
}
});
window.addEventListener('beforeunload', function() {
socket.emit('leave_job', {job_id: jobId});
});
/* --- Cancel button --- */
if (cancelBtn) {
cancelBtn.addEventListener('click', function() {
-1
View File
@@ -35,6 +35,5 @@
</div>
</footer>
{% block scripts %}{% endblock %}
<script src="https://cdn.socket.io/4.7.5/socket.io.min.js"></script>
</body>
</html>
+1 -1
View File
@@ -35,7 +35,7 @@
<div class="live-preview-section card" id="live-preview-section" style="display: block">
<h3 class="section-title">Live Preview</h3>
<div class="live-preview-wrap">
<img id="live-preview-img" class="live-preview-img" alt="Annotated frame from video processing" />
<img id="live-preview-img" class="live-preview-img" src="/api/preview/{{ job.job_id }}" alt="Annotated frame from video processing" />
<div class="preview-placeholder" id="preview-placeholder">Waiting for first frame...</div>
</div>
</div>