#include "voice_interaction_internal.h" #include #include #include #include #include "cJSON.h" #include "esp_audio_dec.h" #include "esp_audio_enc.h" #include "esp_log.h" #include "voice_audio.h" static const char *TAG = "voice_interaction"; static char *voice_build_directive_message(const char *action, const char *directive, bool with_dialog_id) { cJSON *root = cJSON_CreateObject(); if (root == NULL) { return NULL; } cJSON *header = cJSON_AddObjectToObject(root, "header"); cJSON_AddStringToObject(header, "action", action); cJSON_AddStringToObject(header, "task_id", s_voice.task_id); cJSON_AddStringToObject(header, "streaming", VOICE_STREAMING_MODE); cJSON *payload = cJSON_AddObjectToObject(root, "payload"); cJSON *input = cJSON_AddObjectToObject(payload, "input"); cJSON_AddStringToObject(input, "directive", directive); if (with_dialog_id && s_voice.dialog_id[0] != '\0') { cJSON_AddStringToObject(input, "dialog_id", s_voice.dialog_id); } char *text = cJSON_PrintUnformatted(root); cJSON_Delete(root); return text; } static char *voice_build_start_message(void) { cJSON *root = cJSON_CreateObject(); if (root == NULL) { return NULL; } cJSON *header = cJSON_AddObjectToObject(root, "header"); cJSON_AddStringToObject(header, "action", "run-task"); cJSON_AddStringToObject(header, "task_id", s_voice.task_id); cJSON_AddStringToObject(header, "streaming", VOICE_STREAMING_MODE); cJSON *payload = cJSON_AddObjectToObject(root, "payload"); cJSON_AddStringToObject(payload, "task_group", VOICE_TASK_GROUP); cJSON_AddStringToObject(payload, "task", VOICE_TASK_NAME); cJSON_AddStringToObject(payload, "function", VOICE_FUNCTION_NAME); cJSON_AddStringToObject(payload, "model", VOICE_MODEL_NAME); cJSON *input = cJSON_AddObjectToObject(payload, "input"); cJSON_AddStringToObject(input, "directive", "Start"); cJSON_AddStringToObject(input, "workspace_id", CONFIG_TQ_VOICE_WORKSPACE_ID); cJSON_AddStringToObject(input, "app_id", CONFIG_TQ_VOICE_APP_ID); cJSON *parameters = cJSON_AddObjectToObject(payload, "parameters"); cJSON *upstream = cJSON_AddObjectToObject(parameters, "upstream"); cJSON_AddStringToObject(upstream, "type", "AudioOnly"); cJSON_AddStringToObject(upstream, "mode", "tap2talk"); cJSON_AddStringToObject(upstream, "audio_format", "raw-opus"); cJSON_AddNumberToObject(upstream, "sample_rate", CONFIG_TQ_VOICE_SAMPLE_RATE); cJSON *downstream = cJSON_AddObjectToObject(parameters, "downstream"); cJSON_AddNumberToObject(downstream, "sample_rate", CONFIG_TQ_VOICE_SAMPLE_RATE); cJSON_AddStringToObject(downstream, "audio_format", "raw-opus"); cJSON_AddNumberToObject(downstream, "frame_size", CONFIG_TQ_VOICE_OPUS_FRAME_MS); cJSON_AddNumberToObject(downstream, "bit_rate", CONFIG_TQ_VOICE_OPUS_BITRATE_KBPS); if (strlen(CONFIG_TQ_VOICE_TTS_VOICE) > 0) { cJSON_AddStringToObject(downstream, "voice", CONFIG_TQ_VOICE_TTS_VOICE); } cJSON *client_info = cJSON_AddObjectToObject(parameters, "client_info"); cJSON_AddStringToObject(client_info, "user_id", voice_safe_user_id()); cJSON *device = cJSON_AddObjectToObject(client_info, "device"); cJSON_AddStringToObject(device, "uuid", s_voice.device_uuid); char *text = cJSON_PrintUnformatted(root); cJSON_Delete(root); return text; } static esp_err_t voice_send_text(const char *text, int len) { if (text == NULL || len <= 0) { return ESP_ERR_INVALID_ARG; } esp_websocket_client_handle_t ws = NULL; if (!voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { return ESP_ERR_TIMEOUT; } ws = s_voice.ws; bool connected = s_voice.ws_connected; voice_unlock(); if (ws == NULL || !connected) { ESP_LOGW(TAG, "[stage] ws_send_text skipped: ws=%p connected=%d len=%d", (void *)ws, connected, len); return ESP_ERR_INVALID_STATE; } int ret = esp_websocket_client_send_text(ws, text, len, pdMS_TO_TICKS(VOICE_WS_SEND_TIMEOUT_MS)); if (ret < 0) { ESP_LOGW(TAG, "[stage] ws_send_text failed: len=%d", len); return ESP_FAIL; } ESP_LOGI(TAG, "[stage] ws_send_text ok: len=%d", len); return ESP_OK; } esp_err_t voice_send_directive(const char *action, const char *directive, bool with_dialog_id) { char *text = NULL; if (!voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { return ESP_ERR_TIMEOUT; } text = voice_build_directive_message(action, directive, with_dialog_id); voice_unlock(); if (text == NULL) { return ESP_ERR_NO_MEM; } esp_err_t err = voice_send_text(text, (int)strlen(text)); ESP_LOGI(TAG, "[stage] directive: action=%s directive=%s with_dialog_id=%d result=%s", action, directive, with_dialog_id, esp_err_to_name(err)); cJSON_free(text); return err; } static void voice_handle_downstream_packet(const uint8_t *data, size_t len) { if (data == NULL || len == 0) { return; } if (!voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { return; } void *opus_dec = s_voice.opus_dec; uint8_t *pcm_buf = s_voice.pcm_rx_buf; size_t pcm_buf_size = s_voice.pcm_rx_buf_size; voice_unlock(); if (opus_dec == NULL || pcm_buf == NULL || pcm_buf_size == 0) { return; } esp_audio_dec_in_raw_t raw = { .buffer = (uint8_t *)data, .len = (uint32_t)len, .consumed = 0, }; while (raw.len > 0) { esp_audio_dec_out_frame_t out = { .buffer = pcm_buf, .len = (uint32_t)pcm_buf_size, .needed_size = 0, .decoded_size = 0, }; esp_audio_dec_info_t dec_info = {0}; esp_audio_err_t ret = esp_opus_dec_decode(opus_dec, &raw, &out, &dec_info); if (ret == ESP_AUDIO_ERR_BUFF_NOT_ENOUGH) { if (out.needed_size > pcm_buf_size) { uint8_t *new_buf = (uint8_t *)voice_realloc_prefer_psram(pcm_buf, out.needed_size); if (new_buf == NULL) { ESP_LOGE(TAG, "No memory to extend pcm rx buffer to %u", out.needed_size); break; } if (voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { s_voice.pcm_rx_buf = new_buf; s_voice.pcm_rx_buf_size = out.needed_size; voice_unlock(); } pcm_buf = new_buf; pcm_buf_size = out.needed_size; continue; } break; } if (ret != ESP_AUDIO_ERR_OK) { ESP_LOGW(TAG, "opus decode failed: %d", ret); break; } if (out.decoded_size > 0) { char audio_err[64] = {0}; size_t samples = (size_t)out.decoded_size / VOICE_SAMPLE_BYTES; esp_err_t write_err = voice_audio_write_pcm((const int16_t *)out.buffer, samples, VOICE_AUDIO_IO_TIMEOUT_MS, audio_err, sizeof(audio_err)); if (write_err == ESP_OK && voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { int64_t now_ms = voice_now_ms(); int64_t pcm_ms = ((int64_t)samples * 1000 + (int64_t)CONFIG_TQ_VOICE_SAMPLE_RATE - 1) / (int64_t)CONFIG_TQ_VOICE_SAMPLE_RATE; int64_t base_ms = s_voice.playback_deadline_ms > now_ms ? s_voice.playback_deadline_ms : now_ms; s_voice.playback_deadline_ms = base_ms + pcm_ms + VOICE_AUDIO_DRAIN_MARGIN_MS; s_voice.last_downstream_ms = now_ms; voice_unlock(); } } if (raw.consumed == 0 || raw.consumed > raw.len) { break; } raw.buffer += raw.consumed; raw.len -= raw.consumed; } if (voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { s_voice.downstream_packets++; s_voice.last_event_ms = voice_now_ms(); voice_unlock(); } } static const char *voice_get_json_str(cJSON *obj, const char *name) { if (obj == NULL || name == NULL) { return NULL; } cJSON *item = cJSON_GetObjectItemCaseSensitive(obj, name); if (cJSON_IsString(item) && item->valuestring != NULL) { return item->valuestring; } return NULL; } static bool voice_get_json_bool(cJSON *obj, const char *name, bool *out_value) { if (obj == NULL || name == NULL || out_value == NULL) { return false; } cJSON *item = cJSON_GetObjectItemCaseSensitive(obj, name); if (!cJSON_IsBool(item)) { return false; } *out_value = cJSON_IsTrue(item); return true; } static void voice_log_final_text(const char *stage, const char *text) { if (stage == NULL || text == NULL || text[0] == '\0') { return; } size_t len = strlen(text); const int preview = 240; ESP_LOGI(TAG, "[stage] %s: len=%u text=%.*s%s", stage, (unsigned)len, preview, text, (len > (size_t)preview) ? "..." : ""); } static void voice_handle_output_event(cJSON *output) { const char *event_name = voice_get_json_str(output, "event"); if (event_name == NULL) { return; } const char *dialog_id = voice_get_json_str(output, "dialog_id"); const char *state = NULL; if (strcmp(event_name, "DialogStateChanged") == 0) { state = voice_get_json_str(output, "state"); } ESP_LOGI(TAG, "[stage] ws_event_output: event=%s dialog_id=%s state=%s", event_name, dialog_id != NULL ? dialog_id : "-", state != NULL ? state : "-"); bool finished = false; bool has_finished = voice_get_json_bool(output, "finished", &finished); if (has_finished && finished) { if (strcmp(event_name, "SpeechContent") == 0) { voice_log_final_text("asr_final_text", voice_get_json_str(output, "text")); } else if (strcmp(event_name, "RespondingContent") == 0) { const char *final_text = voice_get_json_str(output, "text"); if (final_text == NULL || final_text[0] == '\0') { final_text = voice_get_json_str(output, "spoken"); } voice_log_final_text("response_final_text", final_text); } } bool should_close_audio = false; if (voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { if (dialog_id != NULL) { strlcpy(s_voice.dialog_id, dialog_id, sizeof(s_voice.dialog_id)); } s_voice.last_event_ms = voice_now_ms(); if (strcmp(event_name, "Started") == 0) { s_voice.started = true; s_voice.dialog_state = VOICE_DIALOG_STATE_IDLE; ESP_LOGI(TAG, "[stage] session_started: dialog_id=%s", s_voice.dialog_id); } else if (strcmp(event_name, "DialogStateChanged") == 0) { if (state != NULL) { if (strcmp(state, "Listening") == 0) { s_voice.dialog_state = VOICE_DIALOG_STATE_LISTENING; } else if (strcmp(state, "Thinking") == 0) { s_voice.dialog_state = VOICE_DIALOG_STATE_THINKING; } else if (strcmp(state, "Responding") == 0) { s_voice.dialog_state = VOICE_DIALOG_STATE_RESPONDING; } } ESP_LOGI(TAG, "[stage] dialog_state_changed: state=%s current=%s", state != NULL ? state : "-", voice_interaction_dialog_state_str(s_voice.dialog_state)); } else if (strcmp(event_name, "SpeechEnded") == 0) { ESP_LOGI(TAG, "[stage] speech_ended: keep tap_active=%d for continuous rounds", s_voice.tap_active); } else if (strcmp(event_name, "Stopped") == 0) { s_voice.started = false; s_voice.tap_active = false; s_voice.dialog_state = VOICE_DIALOG_STATE_IDLE; should_close_audio = true; ESP_LOGI(TAG, "[stage] session_stopped_by_server"); } else if (strcmp(event_name, "Error") == 0) { const char *error_msg = voice_get_json_str(output, "error_message"); if (error_msg != NULL) { voice_set_last_error_locked(error_msg); ESP_LOGW(TAG, "[stage] server_error: %s", error_msg); } s_voice.dialog_state = VOICE_DIALOG_STATE_IDLE; s_voice.tap_active = false; should_close_audio = true; } voice_unlock(); } if (should_close_audio) { voice_close_audio_with_drain(VOICE_AUDIO_DRAIN_WAIT_MS, "server_event_close"); } if (strcmp(event_name, "RespondingStarted") == 0) { (void)voice_send_directive("continue-task", "LocalRespondingStarted", true); } else if (strcmp(event_name, "RespondingEnded") == 0) { (void)voice_send_directive("continue-task", "LocalRespondingEnded", true); } } static void voice_handle_text_message(const char *text, size_t len) { cJSON *root = cJSON_ParseWithLength(text, len); if (root == NULL) { ESP_LOGW(TAG, "Invalid ws text payload"); return; } cJSON *header = cJSON_GetObjectItemCaseSensitive(root, "header"); const char *header_event = voice_get_json_str(header, "event"); cJSON *payload = cJSON_GetObjectItemCaseSensitive(root, "payload"); cJSON *output = NULL; if (payload != NULL) { output = cJSON_GetObjectItemCaseSensitive(payload, "output"); } if (header != NULL) { cJSON *status_code = cJSON_GetObjectItemCaseSensitive(header, "status_code"); cJSON *status_msg = cJSON_GetObjectItemCaseSensitive(header, "status_message"); if (cJSON_IsNumber(status_code) && status_code->valueint >= 400) { ESP_LOGW(TAG, "[stage] ws_header_error: status_code=%d status_message=%s", status_code->valueint, (cJSON_IsString(status_msg) && status_msg->valuestring != NULL) ? status_msg->valuestring : "-"); if (voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { if (cJSON_IsString(status_msg) && status_msg->valuestring != NULL) { voice_set_last_error_locked(status_msg->valuestring); } else { voice_set_last_error_locked("websocket status error"); } voice_unlock(); } } } if (output != NULL) { voice_handle_output_event(output); } else if (header_event != NULL && strcmp(header_event, "task-failed") == 0) { const char *error_msg = voice_get_json_str(header, "error_message"); ESP_LOGW(TAG, "[stage] ws_task_failed: %s", error_msg != NULL ? error_msg : "task failed"); if (voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { voice_set_last_error_locked(error_msg != NULL ? error_msg : "task failed"); voice_unlock(); } } cJSON_Delete(root); } static void voice_handle_text_chunk(esp_websocket_event_data_t *data) { if (data->payload_offset == 0) { free(s_voice.text_agg); s_voice.text_agg = NULL; s_voice.text_agg_size = 0; if (data->payload_len <= 0 || data->payload_len > 4096) { return; } s_voice.text_agg = (uint8_t *)voice_calloc_prefer_psram(1, (size_t)data->payload_len + 1); if (s_voice.text_agg == NULL) { ESP_LOGW(TAG, "[stage] text_agg_alloc_failed: len=%d", data->payload_len); return; } s_voice.text_agg_size = (size_t)data->payload_len; } if (s_voice.text_agg == NULL || s_voice.text_agg_size == 0) { return; } if ((size_t)data->payload_offset + (size_t)data->data_len > s_voice.text_agg_size) { free(s_voice.text_agg); s_voice.text_agg = NULL; s_voice.text_agg_size = 0; return; } memcpy(s_voice.text_agg + data->payload_offset, data->data_ptr, (size_t)data->data_len); bool complete = data->fin && ((size_t)data->payload_offset + (size_t)data->data_len == s_voice.text_agg_size); if (!complete) { return; } voice_handle_text_message((const char *)s_voice.text_agg, s_voice.text_agg_size); free(s_voice.text_agg); s_voice.text_agg = NULL; s_voice.text_agg_size = 0; } static void voice_handle_binary_chunk(esp_websocket_event_data_t *data) { if (data->payload_offset == 0) { free(s_voice.bin_agg); s_voice.bin_agg = NULL; s_voice.bin_agg_size = 0; if (data->payload_len <= 0 || data->payload_len > 16384) { return; } s_voice.bin_agg = (uint8_t *)voice_malloc_prefer_psram((size_t)data->payload_len); if (s_voice.bin_agg == NULL) { ESP_LOGW(TAG, "[stage] bin_agg_alloc_failed: len=%d", data->payload_len); return; } s_voice.bin_agg_size = (size_t)data->payload_len; } if (s_voice.bin_agg == NULL || s_voice.bin_agg_size == 0) { return; } if ((size_t)data->payload_offset + (size_t)data->data_len > s_voice.bin_agg_size) { free(s_voice.bin_agg); s_voice.bin_agg = NULL; s_voice.bin_agg_size = 0; return; } memcpy(s_voice.bin_agg + data->payload_offset, data->data_ptr, (size_t)data->data_len); bool complete = data->fin && ((size_t)data->payload_offset + (size_t)data->data_len == s_voice.bin_agg_size); if (!complete) { return; } voice_handle_downstream_packet(s_voice.bin_agg, s_voice.bin_agg_size); free(s_voice.bin_agg); s_voice.bin_agg = NULL; s_voice.bin_agg_size = 0; } void voice_websocket_event_handler(void *handler_args, esp_event_base_t base, int32_t event_id, void *event_data) { (void)handler_args; (void)base; UBaseType_t ws_hwm = uxTaskGetStackHighWaterMark(NULL); if (!s_voice.ws_low_stack_warned && ws_hwm < 256) { s_voice.ws_low_stack_warned = true; ESP_LOGW(TAG, "[stage] ws_task_low_stack: hwm_words=%u", (unsigned)ws_hwm); } esp_websocket_event_data_t *data = (esp_websocket_event_data_t *)event_data; switch ((esp_websocket_event_id_t)event_id) { case WEBSOCKET_EVENT_CONNECTED: { ESP_LOGI(TAG, "[stage] ws_connected"); if (voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { s_voice.ws_connected = true; s_voice.last_event_ms = voice_now_ms(); voice_unlock(); } char *start_msg = NULL; if (voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { start_msg = voice_build_start_message(); voice_unlock(); } if (start_msg != NULL) { esp_err_t send_err = voice_send_text(start_msg, (int)strlen(start_msg)); ESP_LOGI(TAG, "[stage] start_message_sent: result=%s", esp_err_to_name(send_err)); cJSON_free(start_msg); } else { ESP_LOGE(TAG, "[stage] start_message_build_failed"); } break; } case WEBSOCKET_EVENT_DISCONNECTED: ESP_LOGW(TAG, "[stage] ws_disconnected"); if (voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { s_voice.ws_connected = false; s_voice.started = false; s_voice.tap_active = false; s_voice.dialog_state = VOICE_DIALOG_STATE_IDLE; s_voice.last_event_ms = voice_now_ms(); voice_unlock(); } voice_close_audio_with_drain(VOICE_AUDIO_DRAIN_WAIT_MS, "ws_disconnected_close"); break; case WEBSOCKET_EVENT_DATA: if (data == NULL || data->data_ptr == NULL || data->data_len <= 0) { break; } if (data->op_code == WS_TRANSPORT_OPCODES_TEXT) { voice_handle_text_chunk(data); } else if (data->op_code == WS_TRANSPORT_OPCODES_BINARY) { voice_handle_binary_chunk(data); } break; case WEBSOCKET_EVENT_ERROR: ESP_LOGW(TAG, "[stage] ws_error"); if (voice_lock(VOICE_STATUS_LOCK_TIMEOUT_MS)) { voice_set_last_error_locked("websocket error"); s_voice.last_event_ms = voice_now_ms(); voice_unlock(); } break; default: break; } }