Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .gitlab-ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -89,13 +89,13 @@ test:memory-leak:
# Leak-check the pure (no-subprocess) binaries, grouped by module:
# rtc-tcp-client : tai_unit_tests, tai_log_level_zero_tests,
# tai_log_module_zero_tests, tai_integration_tests
# iot-client : iot_cipher_test
# iot-client : iot_cipher_test, iot_ai_ctrl_test
# tuya-ble : tuya_ble_test
# rtc-client is a prebuilt closed lib (headers + libstm.a) with no offline
# test, so there is nothing to leak-check for it here. The mock-driven iot
# tests are covered by the `test` job; under valgrind their Python
# handshakes would time out.
for t in tai_unit_tests tai_log_level_zero_tests tai_log_module_zero_tests tai_integration_tests iot_cipher_test tuya_ble_test; do
for t in tai_unit_tests tai_log_level_zero_tests tai_log_module_zero_tests tai_integration_tests iot_cipher_test iot_ai_ctrl_test tuya_ble_test; do
echo "==================== valgrind: ${t} ===================="
"$VG" "./${CMAKE_BUILD_DIR}/${t}" 2>&1 | tee "valgrind-${t}.log"
rc=${PIPESTATUS[0]}
Expand Down
15 changes: 15 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,21 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Added

- iot-client — MQTT protocol-9000 AI control channel for server-initiated `asrInterrupt` notices, delivered via `ai_ctrl_callback_t` independently of the RTC TCP Connection (#42).
- The notice payload carries `eventId` and the server time `time`; the time is the interruption cutoff.
- rtc-tcp-client — `on_flow_control` admission hook for TCP receive backpressure; pauses all inbound Frames when the application's audio queue is full, including codec-frame callbacks mid-Packet (#42).
- New knobs: `AGENTIC_KIT_TAI_FLOW_CONTROL_POLL_MS` (admission-check period while paused, default 50 ms) and `AGENTIC_KIT_TAI_WORKER_YIELD_MS` (worker drain-pass yield under sustained traffic, default 10 ms).
- rtc-tcp-client — `tai_event_msg_t` gains borrowed `user_data` / `user_data_len` (attr 111), so an interruption notice can be read without the SDK parsing JSON (#42).
- `TAI_EVT_CHAT_BREAK` carries `{"breakAttributes":{"time":"<server-time>"}}`; compare that time with the `timestamp_ms` latched from audio START to discard an interrupted stream. A notice whose time is missing or unusable fails closed onto the in-flight stream.

### Changed

- **BREAKING** pal — `pal_t` gains a mandatory `sleep_ms` member; custom PALs must supply it and all consumers must rebuild (#42).
- rtc-tcp-client — audio Packets paused mid-body now retain a pending-delivery cursor instead of silently dropping codec frames; reopening admission resumes from the first unadmitted byte (#42).
- A pinned Packet whose recorded wire length cannot be consumed exactly once fails fast with `TAI_PROTO_ERR_FRAME_DECODE` instead of sliding the receive buffer out of range.
- A header-only START/ONE_SHOT now reaches `on_audio` with `len == 0`, so its server-side `timestamp_ms` can be latched for interruption filtering.
- rtc-tcp-client — `tai_connect()` returns once the session acknowledgement is handled; media the server coalesced into the handshake stays buffered and is first delivered by the receive worker after connect returns (#42).

- docs-site — English edition of the full documentation site, published at `/en/` with Simplified Chinese retained at `/`.
- All 29 docs are mirrored under `docs-site/i18n/en/`, with a locale selector in both the landing-page navbar and the custom docs topbar, localized navbar/footer/sidebar catalogs, and English SVG schematics under `current/images/`.
- Heading anchors use explicit IDs shared across locales, and `npm run check:i18n` gates path, image, and anchor parity as part of `npm run build`.
Expand Down
33 changes: 33 additions & 0 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -193,6 +193,7 @@ set(IOT_CLIENT_SOURCES
"${IOT_CLIENT_DIR}/src/atop_base.c"
"${IOT_CLIENT_DIR}/src/cipher_wrapper.c"
"${IOT_CLIENT_DIR}/src/http_client_interface.c"
"${IOT_CLIENT_DIR}/src/iot_ai_ctrl.c"
"${IOT_CLIENT_DIR}/src/iot_atop.c"
"${IOT_CLIENT_DIR}/src/iot_client.c"
"${IOT_CLIENT_DIR}/src/iot_client_message.c"
Expand Down Expand Up @@ -437,11 +438,43 @@ if(AGENTIC_KIT_BUILD_TESTS AND AGENTIC_KIT_ENABLE_PROJECT_TESTS)
agentic_kit_add_iot_test(iot_reset_test "${IOT_CLIENT_TESTS_DIR}/iot_reset_test.c")
agentic_kit_add_iot_test(iot_mqtt_test "${IOT_CLIENT_TESTS_DIR}/mqtt_test.c")
agentic_kit_add_iot_test(iot_message_test "${IOT_CLIENT_TESTS_DIR}/iot_client_message_test.c")
agentic_kit_add_iot_test(iot_ai_ctrl_test "${IOT_CLIENT_TESTS_DIR}/iot_ai_ctrl_test.c")
agentic_kit_add_iot_test(iot_on_boarding_test "${IOT_CLIENT_TESTS_DIR}/on_boarding_test.c")
agentic_kit_add_iot_test(iot_dp_test "${IOT_CLIENT_TESTS_DIR}/iot_dp_test.c")
agentic_kit_add_iot_test(iot_ota_test "${IOT_CLIENT_TESTS_DIR}/iot_ota_test.c")
agentic_kit_add_iot_test(iot_ota_verify_test "${IOT_CLIENT_TESTS_DIR}/iot_ota_verify_test.c")

add_executable(iot_tai_control_test
${AI_TCP_SOURCES}
"${AI_TCP_TESTS_DIR}/tai_pal_loopback.c"
"${IOT_CLIENT_TESTS_DIR}/iot_tai_control_test.c"
"${IOT_CLIENT_TESTS_DIR}/mqtt_interrupt_demo_test.c"
"${IOT_CLIENT_TESTS_DIR}/test_log.c"
)
target_include_directories(iot_tai_control_test PRIVATE
"${AI_TCP_DIR}/include"
"${AI_TCP_DIR}/src"
"${AI_TCP_TESTS_DIR}"
"${IOT_CLIENT_DIR}/include"
"${IOT_CLIENT_DIR}/src"
"${IOT_CLIENT_TESTS_DIR}"
"${IOT_CLIENT_TESTS_DIR}/log_config"
"${cjson_source_dir}"
"${PAL_DIR}"
"${COMMON_DIR}"
)
target_link_libraries(iot_tai_control_test PRIVATE
tuya_iot_client_testlog
agentic_kit_pal
Threads::Threads
)
target_compile_definitions(iot_tai_control_test PRIVATE
PYTHON3_EXEC="${Python3_EXECUTABLE}"
TEST_CONFIG_DIR="${IOT_CLIENT_TESTS_DIR}/config"
MESSAGE_MOCK_PATH="${IOT_CLIENT_TESTS_DIR}/mock/message_mock.py"
)
add_test(NAME iot_tai_control_test COMMAND iot_tai_control_test)

# The ble sources compile into the executable (not the tuya_ble
# library) so they pick up the test sink's log remap from
# tuya_iot_client_testlog's include path -- that is what lets
Expand Down
6 changes: 6 additions & 0 deletions CONTEXT-MAP.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,12 @@ and decisions.

- **IoT Client → RTC TCP Client**: the IoT Client obtains an AI session token
(`iot_client_get_session_token`) that the RTC client uses to open a Session.
Applications may also keep MQTT active as an independent control path:
authenticated protocol-9000 notices reach `ai_ctrl_callback_t` even while
receive backpressure stalls the RTC Connection. The application composes the
callbacks; neither client owns or calls the other. Decision records:
`modules/iot-client/docs/adr/0002-mqtt-ai-control-remains-application-composed.md`
and `modules/rtc-tcp-client/docs/adr/0001-receive-backpressure-pauses-the-connection.md`.
- **Tuya BLE → IoT Client**: BLE provisioning runs first and hands the device its WiFi
credentials (incl. the pairing token); once on WiFi, the IoT Client uses them to activate
against the cloud.
Expand Down
30 changes: 26 additions & 4 deletions docs-site/docs/guides/vad-and-interrupt.md
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,12 @@ tai_config_t cfg = {

### RTC TCP Client {#rtc-tcp-client}

:::note MQTT 独立控制路径
设备可同时保持 IoT MQTT 与 RTC TCP Connection。注册 `iot_ai_ctrl_set_callback()` 后,protocol-9000 `asrInterrupt` 不受 RTC TCP 接收背压影响。应用应让一个线程独占 `iot_client_process()` / publish / reconnect,并把 MQTT callback 与 `TAI_EVT_CHAT_BREAK` 汇入同一线程安全的播放策略。

打断的判定依据是**服务端时间**,不是 eventId:MQTT `asrInterrupt` 读 `data.time`,TCP `TAI_EVT_CHAT_BREAK` 读 attr 111 中的 `breakAttributes.time`(不是事件负载)。把最大的时间记为截止时间,清空播放队列,再恢复 RTC 接收,让 worker 继续校验和排空旧媒体;不要直接丢弃 TCP 字节,否则会破坏 Frame 边界。`on_audio` 中锁定 START 的 `msg->timestamp_ms`,凡不晚于截止时间的流一律丢弃。服务端打断不要求调用 `tai_chat_break()`,也不要结束或重开云端 VAD 上行流。
:::

**接收服务端打断(`TAI_EVT_CHAT_BREAK`,type=4):**

`TAI_EVT_CHAT_BREAK` 有双重身份:用户在 AI 回复中插话时它是打断信号;在云端 VAD 模式下它同时也是**回合结束信号**(云端检测到用户停止说话后下发,当前云端不再下发 `TAI_EVT_SERVER_VAD`)。两种情况下的设备处理相同:
Expand All @@ -152,18 +158,34 @@ void on_event(tai_ctx_t *ctx, const tai_event_msg_t *msg, void *ud)
if (msg->event_type == TAI_EVT_CHAT_BREAK) {
// 1. 停止 TTS 播放
audio_player_stop();
// 2. 清空播放缓冲区(丢弃本轮在途 TTS,直到下一个 TAI_STREAM_START)
// 2. 从 attr 111 取服务端打断时间,推进截止时间(取最大值)
uint64_t cutoff = parse_break_time(msg->user_data, msg->user_data_len);
if (cutoff) audio_cutoff_ms = cutoff > audio_cutoff_ms ? cutoff : audio_cutoff_ms;
// 3. 清空播放缓冲区;只有晚于截止时间的流才会重新入队
audio_buffer_flush();
// 3. 忽略本轮后续回调
set_ignore_current_response(true);
// 不要停止麦克风采音,不要调用 tai_send_audio_end(),
// 也不要调用 tai_send_audio_start() 重开上行流——
// 云端 VAD 模式下上行 Event 一直保持打开。
// 若该打断的 eventId 没有对应的本地下行缓存,记录日志并忽略即可。
// 若 attr 111 缺失或时间无效,按失败关闭处理:把当前流的 START
// 当作截止时间,丢弃在途流,而不是当作无事发生继续播放。
}
}
```

`on_audio` 侧锁定 START 的服务端时间,并据此过滤:

```c
void on_audio(tai_ctx_t *ctx, const tai_audio_msg_t *msg, void *ud)
{
if (msg->stream_flag == TAI_STREAM_START ||
msg->stream_flag == TAI_STREAM_ONE_SHOT)
stream_start_ms = msg->timestamp_ms; // 服务端时间,非本地时间
// MIDDLE/END 不带新时间,沿用本轮流锁定的起始时间
if (msg->len && stream_start_ms > audio_cutoff_ms)
audio_buffer_push(msg->data, msg->len);
}
```

**发送客户端打断:**

```c
Expand Down
14 changes: 14 additions & 0 deletions docs-site/docs/reference/iot-client.md
Original file line number Diff line number Diff line change
Expand Up @@ -441,6 +441,20 @@ int iot_client_publish(iot_client_t *client, const uint8_t *data, size_t data_le

---

### `iot_ai_ctrl_set_callback` {#iot_ai_ctrl_set_callback}

```c
int iot_ai_ctrl_set_callback(iot_client_t *client,
ai_ctrl_callback_t callback,
void *user_data);
```

注册经过 P2.3 解密和认证的 MQTT protocol-9000 AI 控制通知。回调在调用 `iot_client_process()` 的线程上触发,`type` 和 `json_data` 只在回调期间有效。回调应仅更新有界状态或通知应用线程,不应在其中断开/销毁 IoT client。

若必须接收订阅建立期间立即到达的通知,初始化时设置 `mqtt_disable_auto_connect=true`,先注册回调,再调用 `iot_client_connect()`。传 NULL callback 可注销。RTC TCP 接收背压不会阻塞该 MQTT 路径;应用应把 MQTT 通知和 TAI ChatBreak 汇入同一播放策略。`asrInterrupt` 的负载携带 `eventId` 与服务端时间 `time`(字符串毫秒),后者是打断判定的截止时间;把它与 `on_audio` 中锁定的 `msg->timestamp_ms` 比较来丢弃过期媒体,而不是按 eventId 判断。服务端通知不是 `tai_chat_break()` 的应答,也不应自动结束云端 VAD 上行流。

---

### `iot_get_qrcode_info` {#iot_get_qrcode_info}

```c
Expand Down
23 changes: 21 additions & 2 deletions docs-site/docs/reference/rtc-tcp-client.md
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,7 @@ sidebar_position: 1
| 字段 | 类型 | 说明 |
|------|------|------|
| `ping_interval_ms` | `uint32_t` | Ping 间隔(0 = 默认 60000ms) |
| `ping_timeout_ms` | `uint32_t` | Ping 超时(0 = 默认 90000ms) |
| `ping_timeout_ms` | `uint32_t` | 接收存活超时,任意入站数据均刷新(0 = 默认 90000ms);主动背压暂停期间不计超时,恢复时获得完整的新预算 |
| `connect_timeout_ms` | `uint32_t` | 连接超时(0 = 默认 5000ms)。分别约束 `tai_connect` 的两个串行等待阶段:先是连接建立(TCP 建连 + TLS 握手,共用一份预算),再是服务端 SessionNew 应答。任一阶段超时即判定连接失败,因此 `tai_connect` 最坏耗时约为该值的 2 倍。 |

### 3.7 测试配置 {#37-测试配置}
Expand All @@ -205,6 +205,7 @@ sidebar_position: 1
| `on_image` | function pointer | 图像数据回调(云端生成的图片) |
| `on_event` | function pointer | 事件回调(MCP、打断、VAD 等) |
| `on_disconnect` | function pointer | 断连回调 |
| `on_flow_control` | `int (*)(tai_ctx_t *, void *)` | 可选接收背压钩子;非零允许接收,0 暂停读取与解析,NULL 不启用背压 |
| `user_data` | `void *` | 透传到所有回调 |

**回调签名:**
Expand All @@ -217,8 +218,20 @@ void (*on_text) (tai_ctx_t *ctx, const tai_text_msg_t *msg, void *use
void (*on_image) (tai_ctx_t *ctx, const tai_image_msg_t *msg, void *user_data);
void (*on_event) (tai_ctx_t *ctx, const tai_event_msg_t *msg, void *user_data);
void (*on_disconnect)(tai_ctx_t *ctx, const tai_disconnect_msg_t *msg, void *user_data);
int (*on_flow_control)(tai_ctx_t *ctx, void *user_data);
```

#### 接收背压(`on_flow_control`)

- worker 在每次读取前及完整 Frame 之间调用钩子。符合规范的服务端不会在握手阶段触发它(应答先于媒体下发);仅当服务端把媒体排在应答之前才可能在握手期被调用。钩子必须非阻塞;应用负责同步共享的队列状态。
- 返回 0 会暂停读取和解析,所有入站流量都会停滞,包括 ChatBreak、ASR 文本、Pong 和 EOF 检测。恢复后先处理已缓冲的完整 Frame,再读取;不完整输入回到有上限的阻塞接收。
- 主动暂停挂起接收存活超时;恢复时获得新的完整 `ping_timeout_ms` 预算。Ping 和停止请求仍执行,Ping 发送失败仍断连。
- 背压也在 Audio Packet 内的每个编解码帧回调前检查。中途暂停会以零拷贝方式保留该 Packet 的剩余字节,恢复时先交付剩余帧、再处理后续 Packet。应用应在 `on_audio` 中按服务端时间(`msg->timestamp_ms`)丢弃过期音频,而不是改写 SDK 接收状态。

默认每 `AGENTIC_KIT_TAI_FLOW_CONTROL_POLL_MS`(50 ms)重查背压,并以 `AGENTIC_KIT_TAI_WORKER_YIELD_MS`(10 ms)在持续流量下让出 CPU。

`pal_t` 新增必填的 `sleep_ms` 回调。它必须非忙等地休眠至少请求的毫秒数,0 为 no-op,且不能依赖 socket 可读性。所有自定义 PAL 和使用方必须用新头文件重新编译。

### 接收消息结构体 {#接收消息结构体}

`tai_audio_msg_t`(音频回调):
Expand All @@ -233,7 +246,7 @@ void (*on_disconnect)(tai_ctx_t *ctx, const tai_disconnect_msg_t *msg, void *use
| `stream_flag` | `uint8_t` | `TAI_STREAM_*`(取自媒体头) |
| `data_id` | `uint16_t` | 数据 ID:`AUDIO_DOWN`(2) / `AUDIO_AUX`(7) |
| `event_id` | `const char *` | turn id(借用);无则为 `""` |
| `timestamp_ms` | `uint64_t` | 流起始时间戳(媒体头) |
| `timestamp_ms` | `uint64_t` | 服务端媒体头时间戳,**非**本地时间;过滤流时以 START 的值锁定,与打断时间比较 |

`tai_text_msg_t`(文本回调):

Expand Down Expand Up @@ -272,6 +285,12 @@ void (*on_disconnect)(tai_ctx_t *ctx, const tai_disconnect_msg_t *msg, void *use
| `data` | `const uint8_t *` | 事件负载(通常为 JSON) |
| `len` | `size_t` | 负载字节数 |
| `event_id` | `const char *` | attr 61(借用);无则为 `""` |
| `user_data` | `const uint8_t *` | attr 111(借用),**非** NUL 结尾;无则为 NULL |
| `user_data_len` | `size_t` | `user_data` 字节数;与事件负载分离,SDK 不解析其 JSON |

:::note ChatBreak 打断时间
`TAI_EVT_CHAT_BREAK` 的服务端时间不在事件负载里,而在 attr 111 中:`{"breakAttributes":{"time":"<server-time>"}}`,值为服务端 epoch 毫秒,与 `on_audio` 的 `timestamp_ms` 同钟同单位。读取 `msg->user_data`、解析 `breakAttributes.time`,再与 `on_audio` 中锁定的 `timestamp_ms` 比较,即可判定该流是否已过期。MQTT 侧的 `asrInterrupt` 走另一条路径:时间在其自身负载的 `time` 字段,且是同一个服务端时间值。
:::

`tai_disconnect_msg_t`(断连回调):

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,12 @@ When the user speaks again while the AI is responding, the current response need

### RTC TCP Client {#rtc-tcp-client}

:::note Independent MQTT control path
A device can keep the IoT MQTT and RTC TCP Connection active together. After `iot_ai_ctrl_set_callback()` is registered, a protocol-9000 `asrInterrupt` is independent of RTC TCP receive backpressure. One application thread should own `iot_client_process()`, publish, and reconnect, while the MQTT callback and `TAI_EVT_CHAT_BREAK` feed the same thread-safe playback policy.

An interruption is decided by **server time**, not by event ID: MQTT `asrInterrupt` carries `data.time`, and TCP `TAI_EVT_CHAT_BREAK` carries `breakAttributes.time` inside attr 111 (not in the event payload). Keep the greatest value as the cutoff, flush the playback queue, then release RTC receive pressure so the worker can keep authenticating and draining old media -- never discard raw TCP bytes, which would corrupt Frame boundaries. In `on_audio`, latch START's `msg->timestamp_ms` and drop every stream at or before the cutoff. A server notice does not require `tai_chat_break()`, and it must not end or reopen a Server-VAD uplink.
:::

**Receive a server-initiated chat break (`TAI_EVT_CHAT_BREAK`, type=4):**

`TAI_EVT_CHAT_BREAK` has two roles: it signals an interruption when the user speaks during the AI response; in server-VAD mode, it is also the **turn-end signal** (sent when the cloud detects that the user has stopped speaking; the current cloud no longer sends `TAI_EVT_SERVER_VAD`). Device-side handling is the same in both cases:
Expand All @@ -154,20 +160,34 @@ void on_event(tai_ctx_t *ctx, const tai_event_msg_t *msg, void *ud)
if (msg->event_type == TAI_EVT_CHAT_BREAK) {
// 1. Stop TTS playback
audio_player_stop();
// 2. Clear the playback buffer (discard this turn's in-flight TTS
// until the next TAI_STREAM_START)
// 2. Read the server interruption time from attr 111 and advance the cutoff
uint64_t cutoff = parse_break_time(msg->user_data, msg->user_data_len);
if (cutoff) audio_cutoff_ms = cutoff > audio_cutoff_ms ? cutoff : audio_cutoff_ms;
// 3. Clear the playback buffer; only streams newer than the cutoff requeue
audio_buffer_flush();
// 3. Ignore subsequent callbacks for this turn
set_ignore_current_response(true);
// Do not stop microphone capture or call tai_send_audio_end().
// Do not call tai_send_audio_start() to reopen the uplink stream either:
// in server-VAD mode, the uplink Event remains open throughout.
// If this chat break's eventId has no corresponding local downlink buffer,
// just log and ignore it.
// If attr 111 is missing or its time is invalid, fail closed: treat the
// in-flight stream's START as the cutoff instead of letting it play on.
}
}
```

`on_audio` latches START's server time and filters on it:

```c
void on_audio(tai_ctx_t *ctx, const tai_audio_msg_t *msg, void *ud)
{
if (msg->stream_flag == TAI_STREAM_START ||
msg->stream_flag == TAI_STREAM_ONE_SHOT)
stream_start_ms = msg->timestamp_ms; // server time, not local time
// MIDDLE/END carry no new time; they reuse this turn's latched start
if (msg->len && stream_start_ms > audio_cutoff_ms)
audio_buffer_push(msg->data, msg->len);
}
```

**Send a client-initiated chat break:**

```c
Expand Down
Loading
Loading