From 9c50ad349ee8d616470f80ed19ff2840042b3fa7 Mon Sep 17 00:00:00 2001 From: Adolfo Reyna Date: Fri, 19 Jun 2026 19:53:21 -0400 Subject: [PATCH] Update main.py to use WebSocket audio streaming automatically at boot --- main.py | 246 ++++++++++++++++++++++++++++++-------------------------- 1 file changed, 132 insertions(+), 114 deletions(-) diff --git a/main.py b/main.py index 2ca7cf1..4e61b4b 100644 --- a/main.py +++ b/main.py @@ -35,6 +35,7 @@ from video_stream import VideoStreamServer from audio_util import ES8311 import struct import urequests +from websocket_client import WebSocketClient def create_wav_header(data_size): # Generates a 44-byte WAV header for 16kHz, 16-bit mono PCM @@ -408,158 +409,175 @@ def main(): bits=16, format=I2S.STEREO) - total_data_bytes = 0 - buffer = bytearray(2048) - mono_buf = bytearray(1024) # Half size for mono extraction - - rec_start_time = time.ticks_ms() - max_rec_duration_ms = 10000 # 10 seconds max duration - + ws_connected = False + ws = None try: - with open("voice_rec.pcm", "wb") as f: - while True: - # Safety: check duration + # Determine dynamic device ID based on board configuration + device_id = "esp32_screen" if board_config.BOARD_TYPE == 'WAVESHARE_RLCD' else "little32" + + headers = { + "Authorization": "Bearer mcT1YA1vOr9wXSiHpCYalweEGGZKX-PIfZv2drp8BSg", + "X-Device-ID": device_id + } + ws = WebSocketClient("ws://192.168.68.126:8642/api/esp32/voice/ws", headers=headers) + ws.connect() + ws.send_text(json.dumps({ + "event": "start", + "device_id": device_id, + "sample_rate": 16000, + "channels": 1, + "sample_width": 2, + "format": "pcm_s16le" + })) + ws.recv_frame() # ready + ws.recv_frame() # listening + ws_connected = True + except Exception as wse: + print("WebSocket connect error:", wse) + draw_status_bar("Connection Error") + time.sleep(2) + + if ws_connected and ws: + draw_status_bar("Recording & streaming...") + + # Open I2S RX for recording (Stereo 16kHz) + i2s_rx = I2S(1, + sck=Pin(board_config.audio_i2s_sck), + ws=Pin(board_config.audio_i2s_ws), + sd=Pin(board_config.audio_i2s_rx_sd), + mode=I2S.RX, + ibuf=16000, + rate=16000, + bits=16, + format=I2S.STEREO) + + total_data_bytes = 0 + buffer = bytearray(2048) + mono_buf = bytearray(1024) + + rec_start_time = time.ticks_ms() + max_rec_duration_ms = 10000 # 10 seconds max duration + + try: + # Stream chunks while talk trigger is active + while is_talk_trigger_active(): elapsed = time.ticks_diff(time.ticks_ms(), rec_start_time) if elapsed >= max_rec_duration_ms: print("Recording stopped: maximum duration reached") break - # Read I2S chunk bytes_read = i2s_rx.readinto(buffer) if bytes_read > 0: - # Stereo-to-mono: extract left channel (every other 16-bit sample) mono_len = bytes_read // 2 j = 0 for i in range(0, bytes_read, 4): mono_buf[j] = buffer[i] mono_buf[j + 1] = buffer[i + 1] j += 2 - f.write(mono_buf[:mono_len]) + ws.send_binary(mono_buf[:mono_len]) total_data_bytes += mono_len - except Exception as e: - print("Error recording:", e) - finally: - i2s_rx.deinit() - if board_config.audio_mic_codec == "ES8311": + except Exception as e: + print("Error during recording/streaming:", e) + finally: + i2s_rx.deinit() + if board_config.audio_mic_codec == "ES8311": + try: + codec._write(0x16, 0x00) # Reset mic gain + except: + pass + + if total_data_bytes < 3200: + print("Recording too short, cancelling.") try: - codec._write(0x16, 0x00) # Reset mic gain + ws.send_text(json.dumps({"event": "cancel"})) + ws.close() except: pass - - if total_data_bytes >= 1000: - draw_status_bar("Sending & processing...") - - # Determine dynamic device ID based on board configuration - device_id = "esp32_screen" if board_config.BOARD_TYPE == 'WAVESHARE_RLCD' else "little32" - - # Headers - headers = { - "Authorization": "Bearer mcT1YA1vOr9wXSiHpCYalweEGGZKX-PIfZv2drp8BSg", - "Content-Type": "audio/wav", - "X-Device-ID": device_id, - "X-Hermes-Screen-Device": device_id, - "X-Hermes-Reply-Mode": "ack" - } - - s = None - try: - status_code, resp_headers, s = stream_hermes_request( - "http://192.168.68.126:8642/api/esp32/voice", - headers, - "voice_rec.pcm" - ) - - if status_code == 200: - content_type = resp_headers.get("content-type", "") - content_len = int(resp_headers.get("content-length", 0)) + draw_status_bar("Cancelled") + time.sleep(1) + else: + draw_status_bar("Processing...") + try: + ws.send_text(json.dumps({"event": "stop"})) - if "audio" in content_type and content_len >= 44: - draw_status_bar("Playing response...") + i2s_tx = None + while True: + opcode, payload = ws.recv_frame() + if opcode is None: + break - # Read and parse WAV header (first 44 bytes) - header = socket_read_exactly(s, 44) - if len(header) == 44 and header[0:4] == b'RIFF' and header[8:12] == b'WAVE': - import struct - fmt_idx = header.find(b'fmt ') - if fmt_idx != -1: - fmt_data = header[fmt_idx+8 : fmt_idx+24] - audio_format, channels, sample_rate, byte_rate, block_align, bits = struct.unpack(' 0: - socket_read_exactly(s, bytes_to_skip) - - # Play using board_config parameters + if evt == "transcript": + txt = event_data.get("text", "") + print(f"Heard: {txt}") + draw_status_bar(f"Heard: {txt[:20]}...") + + elif evt == "thinking": + draw_status_bar("Thinking...") + + elif evt == "response_text": + txt = event_data.get("text", "") + print(f"Response: {txt}") + + elif evt == "audio_start": + draw_status_bar("Playing response...") on_val = 0 if board_config.audio_amp_active_level == 0 else 1 - off_val = 1 if board_config.audio_amp_active_level == 0 else 0 amp_pin.value(on_val) # Enable Amp - i2s_format = I2S.MONO if channels == 1 else I2S.STEREO + i2s_format = I2S.MONO # WebSocket audio response is mono i2s_tx = I2S(1, sck=Pin(board_config.audio_i2s_sck), ws=Pin(board_config.audio_i2s_ws), sd=Pin(board_config.audio_i2s_tx_sd), mode=I2S.TX, ibuf=4096, - rate=sample_rate, - bits=bits, + rate=16000, + bits=16, format=i2s_format) - try: - remaining_bytes = content_len - (data_idx + 8) - while remaining_bytes > 0: - chunk_to_read = min(remaining_bytes, 2048) - try: - chunk = s.recv(chunk_to_read) - except OSError: - break - if not chunk: - break - i2s_tx.write(chunk) - remaining_bytes -= len(chunk) - finally: + + elif evt == "audio_end": + if i2s_tx: time.sleep_ms(150) + off_val = 1 if board_config.audio_amp_active_level == 0 else 0 amp_pin.value(off_val) # Disable Amp i2s_tx.deinit() - else: - # Ack or text response - resp_text = "" - if content_len > 0: - try: - resp_text = socket_read_exactly(s, content_len).decode('utf-8', 'ignore') - except: - pass - print("Non-WAV response received:", resp_text) - if "ack" in resp_text.lower() or "ok" in resp_text.lower() or "success" in resp_text.lower(): - draw_status_bar("Success (Ack received)") - else: - draw_status_bar("Response processed.") - time.sleep(2) - else: - content_len = int(resp_headers.get("content-length", 0)) - resp_text = "Unknown error" - if content_len > 0: - try: - resp_text = socket_read_exactly(s, content_len).decode('utf-8', 'ignore') - except: - pass - print("Error response:", resp_text) - draw_status_bar(f"Error: {status_code}") + i2s_tx = None + + elif evt == "done": + break + + elif evt == "error": + msg = event_data.get("message", "Unknown error") + print(f"Server error: {msg}") + draw_status_bar(f"Error: {msg[:20]}") + time.sleep(2) + break + except Exception as e: + print("Error parsing event text:", e) + + elif opcode == 0x2: # Binary frame (Audio WAV chunk) + if i2s_tx: + chunk = payload + if chunk.startswith(b'RIFF') and len(chunk) > 44: + chunk = chunk[44:] + try: + i2s_tx.write(chunk) + except Exception as e: + print("Error writing to speaker:", e) + except Exception as he: + print("Hermes WS query failed:", he) + draw_status_bar("Connection Error") time.sleep(2) - except Exception as he: - print("Hermes HTTP post failed:", he) - sys.print_exception(he) - draw_status_bar("Connection Error") - time.sleep(2) - finally: - if s: + finally: try: - s.close() + ws.close() except: pass - + # Deinit MCLK PWM if mclk_pwm: try: