From b2bcdb5b796a70ad99e4dbc4f896dd22b503bff7 Mon Sep 17 00:00:00 2001 From: Adolfo Reyna Date: Fri, 19 Jun 2026 13:45:18 -0400 Subject: [PATCH] Add high-speed persistent TCP color video stream port (8083) using tight read loop optimization --- lib/rlcd.py | 2 +- video_stream.py | 205 +++++++++++++++++++++++++++++++++++++++++++++++- 2 files changed, 202 insertions(+), 5 deletions(-) diff --git a/lib/rlcd.py b/lib/rlcd.py index 1b56ddd..d9a995d 100644 --- a/lib/rlcd.py +++ b/lib/rlcd.py @@ -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 diff --git a/video_stream.py b/video_stream.py index a78d441..f79de93 100644 --- a/video_stream.py +++ b/video_stream.py @@ -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