mirror of
https://github.com/smittix/intercept.git
synced 2026-07-22 00:08:10 -07:00
Fix WebSocket waterfall "Invalid frame header" by serializing sends
The fft_reader thread was calling ws.send() concurrently with ws.receive() in the main loop. simple-websocket is not thread-safe for simultaneous read/write, corrupting frame headers. Now the reader thread enqueues frames and only the main loop touches the WebSocket. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -1,6 +1,7 @@
|
|||||||
"""WebSocket-based waterfall streaming with I/Q capture and server-side FFT."""
|
"""WebSocket-based waterfall streaming with I/Q capture and server-side FFT."""
|
||||||
|
|
||||||
import json
|
import json
|
||||||
|
import queue
|
||||||
import subprocess
|
import subprocess
|
||||||
import threading
|
import threading
|
||||||
import time
|
import time
|
||||||
@@ -85,9 +86,23 @@ def init_waterfall_websocket(app: Flask):
|
|||||||
reader_thread = None
|
reader_thread = None
|
||||||
stop_event = threading.Event()
|
stop_event = threading.Event()
|
||||||
claimed_device = None
|
claimed_device = None
|
||||||
|
# Queue for outgoing messages — only the main loop touches ws.send()
|
||||||
|
send_queue = queue.Queue(maxsize=120)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
while True:
|
while True:
|
||||||
|
# Drain send queue first (non-blocking)
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
outgoing = send_queue.get_nowait()
|
||||||
|
except queue.Empty:
|
||||||
|
break
|
||||||
|
try:
|
||||||
|
ws.send(outgoing)
|
||||||
|
except Exception:
|
||||||
|
stop_event.set()
|
||||||
|
break
|
||||||
|
|
||||||
try:
|
try:
|
||||||
msg = ws.receive(timeout=0.1)
|
msg = ws.receive(timeout=0.1)
|
||||||
except TimeoutError:
|
except TimeoutError:
|
||||||
@@ -124,6 +139,12 @@ def init_waterfall_websocket(app: Flask):
|
|||||||
app_module.release_sdr_device(claimed_device)
|
app_module.release_sdr_device(claimed_device)
|
||||||
claimed_device = None
|
claimed_device = None
|
||||||
stop_event.clear()
|
stop_event.clear()
|
||||||
|
# Flush stale frames from previous capture
|
||||||
|
while not send_queue.empty():
|
||||||
|
try:
|
||||||
|
send_queue.get_nowait()
|
||||||
|
except queue.Empty:
|
||||||
|
break
|
||||||
|
|
||||||
# Parse config
|
# Parse config
|
||||||
center_freq = float(data.get('center_freq', 100.0))
|
center_freq = float(data.get('center_freq', 100.0))
|
||||||
@@ -229,13 +250,13 @@ def init_waterfall_websocket(app: Flask):
|
|||||||
'sample_rate': sample_rate,
|
'sample_rate': sample_rate,
|
||||||
}))
|
}))
|
||||||
|
|
||||||
# Start reader thread
|
# Start reader thread — puts frames on queue, never calls ws.send()
|
||||||
def fft_reader(
|
def fft_reader(
|
||||||
proc, ws_ref, stop_evt,
|
proc, _send_q, stop_evt,
|
||||||
_fft_size, _avg_count, _fps,
|
_fft_size, _avg_count, _fps,
|
||||||
_start_freq, _end_freq,
|
_start_freq, _end_freq,
|
||||||
):
|
):
|
||||||
"""Read I/Q from subprocess, compute FFT, send binary frames."""
|
"""Read I/Q from subprocess, compute FFT, enqueue binary frames."""
|
||||||
bytes_per_frame = _fft_size * _avg_count * 2
|
bytes_per_frame = _fft_size * _avg_count * 2
|
||||||
frame_interval = 1.0 / _fps
|
frame_interval = 1.0 / _fps
|
||||||
|
|
||||||
@@ -272,9 +293,10 @@ def init_waterfall_websocket(app: Flask):
|
|||||||
)
|
)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
ws_ref.send(frame)
|
_send_q.put_nowait(frame)
|
||||||
except Exception:
|
except queue.Full:
|
||||||
break
|
# Drop frame if main loop can't keep up
|
||||||
|
pass
|
||||||
|
|
||||||
# Pace to target FPS
|
# Pace to target FPS
|
||||||
elapsed = time.monotonic() - frame_start
|
elapsed = time.monotonic() - frame_start
|
||||||
@@ -288,7 +310,7 @@ def init_waterfall_websocket(app: Flask):
|
|||||||
reader_thread = threading.Thread(
|
reader_thread = threading.Thread(
|
||||||
target=fft_reader,
|
target=fft_reader,
|
||||||
args=(
|
args=(
|
||||||
iq_process, ws, stop_event,
|
iq_process, send_queue, stop_event,
|
||||||
fft_size, avg_count, fps,
|
fft_size, avg_count, fps,
|
||||||
start_freq, end_freq,
|
start_freq, end_freq,
|
||||||
),
|
),
|
||||||
|
|||||||
Reference in New Issue
Block a user