Add high-speed persistent TCP color video stream port (8083) using tight read loop optimization
This commit is contained in:
+1
-1
@@ -150,7 +150,7 @@ class RLCD:
|
||||
print("Error drawing BMP on RLCD:", e)
|
||||
return False
|
||||
|
||||
def draw_rgb565(self, x, y, w, h, data):
|
||||
def draw_rgb565(self, x, y, w, h, data, sync_canvas=True):
|
||||
"""Draw raw RGB565 pixel data converted to 1-bit monochrome on the RLCD."""
|
||||
for cy in range(h):
|
||||
screen_y = y + cy
|
||||
|
||||
+201
-4
@@ -5,15 +5,18 @@ import sys
|
||||
import micropython
|
||||
|
||||
class VideoStreamServer:
|
||||
def __init__(self, display, tcp_port=8081, udp_port=8082):
|
||||
def __init__(self, display, tcp_port=8081, udp_port=8082, color_port=8083):
|
||||
self.display = display
|
||||
self.tcp_port = tcp_port
|
||||
self.udp_port = udp_port
|
||||
self.color_port = color_port
|
||||
|
||||
# Sockets
|
||||
self.tcp_server = None
|
||||
self.tcp_client = None
|
||||
self.udp_sock = None
|
||||
self.color_server = None
|
||||
self.color_client = None
|
||||
|
||||
# State
|
||||
self.active = False
|
||||
@@ -31,6 +34,19 @@ class VideoStreamServer:
|
||||
self.current_frame_id = -1
|
||||
self.chunks_received = 0 # Bitmask for chunks 0-14
|
||||
|
||||
# Color buffering (Pre-allocated, 320x240 RGB565 is 153,600 bytes)
|
||||
self.color_buffer = bytearray(153600)
|
||||
self.color_view = memoryview(self.color_buffer)
|
||||
self.color_bytes_received = 0
|
||||
self.color_header = bytearray(16)
|
||||
self.color_header_view = memoryview(self.color_header)
|
||||
self.color_header_received = 0
|
||||
self.color_payload_len = 0
|
||||
self.color_x = 0
|
||||
self.color_y = 0
|
||||
self.color_w = 0
|
||||
self.color_h = 0
|
||||
|
||||
# Performance Stats
|
||||
self.debug = False
|
||||
self.frames_drawn = 0
|
||||
@@ -115,7 +131,7 @@ class VideoStreamServer:
|
||||
|
||||
def start(self):
|
||||
"""Initialize and start the sockets."""
|
||||
# 1. Start TCP Server
|
||||
# 1. Start TCP Server (Mono)
|
||||
try:
|
||||
self.tcp_server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
self.tcp_server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
@@ -126,7 +142,7 @@ class VideoStreamServer:
|
||||
except Exception as e:
|
||||
print(f"Failed to start TCP stream server: {e}")
|
||||
|
||||
# 2. Start UDP Server
|
||||
# 2. Start UDP Server (Mono)
|
||||
try:
|
||||
self.udp_sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
||||
self.udp_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
@@ -135,6 +151,17 @@ class VideoStreamServer:
|
||||
print(f"UDP Stream Server started on port {self.udp_port}")
|
||||
except Exception as e:
|
||||
print(f"Failed to start UDP stream server: {e}")
|
||||
|
||||
# 3. Start Color TCP Server
|
||||
try:
|
||||
self.color_server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
self.color_server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
self.color_server.bind(('', self.color_port))
|
||||
self.color_server.listen(1)
|
||||
self.color_server.setblocking(False)
|
||||
print(f"Color TCP Stream Server started on port {self.color_port}")
|
||||
except Exception as e:
|
||||
print(f"Failed to start Color TCP stream server: {e}")
|
||||
|
||||
def restart_tcp_server(self):
|
||||
"""Re-initializes the TCP stream server socket after a fatal error."""
|
||||
@@ -176,6 +203,27 @@ class VideoStreamServer:
|
||||
except Exception as e:
|
||||
print(f"Failed to restart UDP stream server: {e}")
|
||||
|
||||
def restart_color_server(self):
|
||||
"""Re-initializes the color TCP stream server socket after a fatal error."""
|
||||
print("Restarting Color TCP Stream Server...")
|
||||
self.close_color_client()
|
||||
try:
|
||||
if self.color_server:
|
||||
self.color_server.close()
|
||||
except:
|
||||
pass
|
||||
self.color_server = None
|
||||
time.sleep_ms(100)
|
||||
try:
|
||||
self.color_server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
self.color_server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
self.color_server.bind(('', self.color_port))
|
||||
self.color_server.listen(1)
|
||||
self.color_server.setblocking(False)
|
||||
print(f"Color TCP Stream Server restarted on port {self.color_port}")
|
||||
except Exception as e:
|
||||
print(f"Failed to restart color TCP stream server: {e}")
|
||||
|
||||
def update(self):
|
||||
"""Updates connection states and checks for incoming video data.
|
||||
Call this periodically in the main execution loop.
|
||||
@@ -194,6 +242,7 @@ class VideoStreamServer:
|
||||
print("Video stream timed out. Returning to dashboard.")
|
||||
self.active = False
|
||||
self.tcp_bytes_received = 0
|
||||
self.color_bytes_received = 0
|
||||
|
||||
# 2. Handle incoming UDP packet (if UDP is active)
|
||||
if self.udp_sock:
|
||||
@@ -246,7 +295,7 @@ class VideoStreamServer:
|
||||
self.restart_udp_sock()
|
||||
break
|
||||
|
||||
# 3. Handle TCP stream
|
||||
# 3. Handle TCP stream (Mono)
|
||||
if self.tcp_server:
|
||||
# If no client is connected, try to accept one
|
||||
if self.tcp_client is None:
|
||||
@@ -296,6 +345,126 @@ class VideoStreamServer:
|
||||
print(f"TCP Stream socket error: {e}")
|
||||
self.close_tcp_client()
|
||||
|
||||
# 4. Handle Color TCP stream
|
||||
if self.color_server:
|
||||
if self.color_client is None:
|
||||
try:
|
||||
self.color_client, addr = self.color_server.accept()
|
||||
self.color_client.setblocking(False)
|
||||
self.color_bytes_received = 0
|
||||
self.color_header_received = 0
|
||||
self.color_payload_len = 0
|
||||
self.active = True
|
||||
self.last_packet_time = now
|
||||
print(f"Color TCP Stream client connected from: {addr}")
|
||||
except OSError as e:
|
||||
import errno
|
||||
err = getattr(e, 'errno', None)
|
||||
if err is None and e.args:
|
||||
err = e.args[0]
|
||||
ewouldblock = getattr(errno, 'EWOULDBLOCK', errno.EAGAIN)
|
||||
if err not in (errno.EAGAIN, ewouldblock) and err is not None:
|
||||
print(f"Fatal Color TCP stream server socket error in accept(): {e}. Re-initializing...")
|
||||
self.restart_color_server()
|
||||
|
||||
if self.color_client is not None:
|
||||
try:
|
||||
# A. Read header first if we haven't got it
|
||||
if self.color_header_received < 16:
|
||||
read_start_time = time.ticks_ms()
|
||||
while self.color_header_received < 16:
|
||||
# 100ms timeout for header read to prevent locking up if stalled
|
||||
if time.ticks_diff(time.ticks_ms(), read_start_time) > 100:
|
||||
break
|
||||
|
||||
remaining_header = 16 - self.color_header_received
|
||||
try:
|
||||
n = self.color_client.readinto(self.color_header_view[self.color_header_received : self.color_header_received + remaining_header])
|
||||
if n is not None:
|
||||
if n == 0:
|
||||
print("Color TCP Stream client disconnected during header read.")
|
||||
self.close_color_client()
|
||||
break
|
||||
else:
|
||||
self.color_header_received += n
|
||||
self.last_packet_time = now
|
||||
self.active = True
|
||||
else:
|
||||
time.sleep_ms(1)
|
||||
except OSError as e:
|
||||
import errno
|
||||
err = getattr(e, 'errno', None)
|
||||
if err is None and e.args:
|
||||
err = e.args[0]
|
||||
ewouldblock = getattr(errno, 'EWOULDBLOCK', errno.EAGAIN)
|
||||
if err in (errno.EAGAIN, ewouldblock) or err is None:
|
||||
time.sleep_ms(1)
|
||||
else:
|
||||
raise e
|
||||
|
||||
# Check if header is now fully received
|
||||
if self.color_header_received == 16:
|
||||
magic = self.color_header[:4]
|
||||
if magic != b'RAW\x01':
|
||||
print(f"Invalid stream magic: {magic}. Disconnecting client.")
|
||||
self.close_color_client()
|
||||
else:
|
||||
self.color_x = (self.color_header[4] << 8) | self.color_header[5]
|
||||
self.color_y = (self.color_header[6] << 8) | self.color_header[7]
|
||||
self.color_w = (self.color_header[8] << 8) | self.color_header[9]
|
||||
self.color_h = (self.color_header[10] << 8) | self.color_header[11]
|
||||
self.color_payload_len = (self.color_header[12] << 24) | (self.color_header[13] << 16) | (self.color_header[14] << 8) | self.color_header[15]
|
||||
self.color_bytes_received = 0
|
||||
if self.color_payload_len > len(self.color_buffer):
|
||||
print(f"Payload too large: {self.color_payload_len} bytes. Disconnecting.")
|
||||
self.close_color_client()
|
||||
|
||||
# B. Read payload if header is fully received
|
||||
if self.color_header_received == 16 and self.color_payload_len > 0:
|
||||
read_start_time = time.ticks_ms()
|
||||
while self.color_bytes_received < self.color_payload_len:
|
||||
# 1000ms timeout for payload read to prevent locking up if stalled
|
||||
if time.ticks_diff(time.ticks_ms(), read_start_time) > 1000:
|
||||
print("Timeout reading color payload.")
|
||||
self.close_color_client()
|
||||
break
|
||||
|
||||
remaining_payload = self.color_payload_len - self.color_bytes_received
|
||||
try:
|
||||
n = self.color_client.readinto(self.color_view[self.color_bytes_received : self.color_bytes_received + remaining_payload])
|
||||
if n is not None:
|
||||
if n == 0:
|
||||
print("Color TCP Stream client disconnected during payload read.")
|
||||
self.close_color_client()
|
||||
break
|
||||
else:
|
||||
self.color_bytes_received += n
|
||||
self.last_packet_time = now
|
||||
self.active = True
|
||||
else:
|
||||
time.sleep_ms(1)
|
||||
except OSError as e:
|
||||
import errno
|
||||
err = getattr(e, 'errno', None)
|
||||
if err is None and e.args:
|
||||
err = e.args[0]
|
||||
ewouldblock = getattr(errno, 'EWOULDBLOCK', errno.EAGAIN)
|
||||
if err in (errno.EAGAIN, ewouldblock) or err is None:
|
||||
time.sleep_ms(1)
|
||||
else:
|
||||
raise e
|
||||
|
||||
# If frame payload is fully received, draw it!
|
||||
if self.color_bytes_received == self.color_payload_len:
|
||||
self._draw_color_frame()
|
||||
# Reset header state for the next frame
|
||||
self.color_header_received = 0
|
||||
self.color_payload_len = 0
|
||||
self.color_bytes_received = 0
|
||||
except OSError as e:
|
||||
print(f"Color TCP Stream socket error: {e}")
|
||||
self.close_color_client()
|
||||
|
||||
def _draw_frame(self):
|
||||
"""Sends the frame buffer directly to the SPI display bus."""
|
||||
draw_start = time.ticks_ms()
|
||||
@@ -317,6 +486,22 @@ class VideoStreamServer:
|
||||
self.frames_drawn += 1
|
||||
self.fps_frame_count += 1
|
||||
|
||||
def _draw_color_frame(self):
|
||||
"""Draws the received color frame directly to the display using draw_rgb565."""
|
||||
draw_start = time.ticks_ms()
|
||||
# Draw the frame using fast path and sync_canvas=False
|
||||
self.display.draw_rgb565(
|
||||
self.color_x,
|
||||
self.color_y,
|
||||
self.color_w,
|
||||
self.color_h,
|
||||
self.color_view[:self.color_payload_len],
|
||||
sync_canvas=False
|
||||
)
|
||||
self.last_draw_ms = time.ticks_diff(time.ticks_ms(), draw_start)
|
||||
self.frames_drawn += 1
|
||||
self.fps_frame_count += 1
|
||||
|
||||
def get_stats(self):
|
||||
"""Returns streaming performance counters."""
|
||||
return {
|
||||
@@ -339,3 +524,15 @@ class VideoStreamServer:
|
||||
pass
|
||||
self.tcp_client = None
|
||||
self.tcp_bytes_received = 0
|
||||
|
||||
def close_color_client(self):
|
||||
"""Safely disconnects the color TCP client socket."""
|
||||
if self.color_client:
|
||||
try:
|
||||
self.color_client.close()
|
||||
except:
|
||||
pass
|
||||
self.color_client = None
|
||||
self.color_bytes_received = 0
|
||||
self.color_header_received = 0
|
||||
self.color_payload_len = 0
|
||||
|
||||
Reference in New Issue
Block a user