Skip to content

Feature/tai backpressure2 - #42

Open
sedawwk wants to merge 10 commits into
masterfrom
feature/tai-backpressure2
Open

sedawwk wants to merge 10 commits into
masterfrom
feature/tai-backpressure2

Conversation

@sedawwk

@sedawwk sedawwk commented Sep 24, 2026

Copy link
Copy Markdown
Collaborator

Summary

Keep server-initiated AI interruptions deliverable while TTS playback applies TCP receive backpressure, without dropping codec frames or corrupting Frame boundaries.

  • Add an authenticated MQTT protocol-9000 control channel through iot_ai_ctrl_set_callback(), independent of the RTC TCP Connection.
  • Add optional tai_config_t.on_flow_control admission checks before receives, between Frames, and before codec-frame callbacks. Paused audio resumes from a zero-copy, worker-owned cursor without replaying or losing accepted bytes.
  • Suspend receive-liveness timeout accounting during intentional pauses, preserve buffered Frames before EOF, and use PAL sleeps to avoid busy polling and worker starvation.
  • Expose borrowed event user_data (attr 111) and demonstrate app-owned interruption filtering by server time rather than event ID. MQTT asrInterrupt and TCP ChatBreak advance a shared cutoff; obsolete audio is authenticated and drained without returning to playback.
  • Add mqtt_interrupt_demo, regression coverage, architectural decision records, and Chinese/English API and interruption documentation.

Application Behavior

Backpressure pauses all inbound RTC TCP traffic, including ChatBreak, Pong, text, and EOF detection. Applications must keep an independent MQTT pump running to receive interruptions while media is paused.

Playback synchronization and stale-audio filtering remain application-owned. The demo flushes playback, releases receive pressure, and accepts only streams whose START timestamp is newer than the interruption cutoff. Notices without a usable time fail closed against the in-flight stream. Server interruptions do not end or reopen the Server-VAD uplink.

Compatibility

  • Breaking: pal_t.sleep_ms is mandatory. Custom PAL ports must implement it, and all consumers must rebuild against the updated public structures. POSIX and FreeRTOS implementations are included.
  • on_flow_control == NULL preserves continuous receive behavior.
  • on_audio can now receive a zero-length header-only START/ONE_SHOT to establish the stream timestamp.
  • Event user_data is callback-lifetime, borrowed, and not NUL-terminated.

Validation

Recorded verification in the branch commits:

  • 20/20 CTest suites pass serially.
  • tai_integration_tests: 7,481 assertions, zero failures; iot_tai_control_test: PASS. Both focused suites pass ASan/UBSan.
  • Regressions cover mid-Packet pause/resume, buffered data before EOF, partial reads, liveness across long pauses, encrypted MQTT interruption ordering, authentication rejection, duplicate notices, reconnect, and shutdown under pressure.
  • mqtt_interrupt_demo, text_chat_demo, and audio_chat_demo build; the documentation site builds with i18n parity.
  • Real speech barge-in confirms that the server interruption cutoff falls between the interrupted stream’s START and the next stream’s START.

sedawwk and others added 10 commits September 29, 2026 13:44
…terrupts

Add a standalone MQTT control channel (protocol 9000) so server-initiated
interrupts (asrInterrupt) arrive immediately via MQTT, independent of the
TCP data channel's congestion state. Mirrors TuyaOpen's architecture: data
over TCP, control over MQTT.

- ai_ctrl_callback_t typedef and iot_ai_ctrl_set_callback() API in iot_client.h
- iot_ai_ctrl.c: parse protocol-9000 envelopes, extract data.data.type and
  data.data.data, deliver to registered callback
- iot_client_message.c: dispatch order is now ai_ctrl (9000) → DP (5) → raw
- 9 network-free unit tests covering dispatch, passthrough, dedup, edge cases
The worker calls the app's on_flow_control hook before each recv and
between buffered frames; when it returns 0 the read is skipped so the
lwIP receive window closes and stalls the sender. Buffered frames resume
without requiring new network data. Add worker_yield (zero-event
tcp_poll) so single-core targets (ESP32-C3) cannot starve the IDLE task
and trip the task watchdog under a TTS flood. New knobs:
TAI_FLOW_CONTROL_POLL_MS and TAI_WORKER_YIELD_MS. Integration test
covers per-packet pause, resume without new bytes, and disconnect while
paused.
A paused worker could still be disconnected by the liveness deadline,
EOF could drop buffered complete Frames, a partial frame fell into a
busy poll, and the drain rechecked admission after the read instead of
before it. Pauses now suspend receive liveness (resume grants a fresh
ping_timeout_ms budget), buffered Frames drain before the next read or
EOF check, and partial input returns to bounded blocking receive.

Add mandatory pal_t.sleep_ms for backpressure waits and CPU yields,
replacing the zero-event tcp_poll trick (POSIX/FreeRTOS nanosleep /
vTaskDelay with tick-rounding chunking). pal_is_valid requires it, so
custom PAL ports must supply it and all consumers must rebuild —
pal_t grows.

Worker drain-pass yields require received or buffered input. Idle
receive timeouts return directly to housekeeping and the next blocking
receive, avoiding the extra AGENTIC_KIT_TAI_WORKER_YIELD_MS delay.

New integration coverage: no reads while paused, buffered Frame
delivered before EOF, partial header/body blocking recv, liveness
across a long pause, and PAL sleep behavior on loopback and POSIX.
Verified against 15/15 ctest suites (tai_integration_tests 24 cases
passing).

Co-Authored-By: Claude Code <[email protected]>
Co-Authored-By: OpenAI Codex <[email protected]>
Show one application-owned MQTT pump alongside pressured TAI media so server interrupts can flush stale playback without corrupting the TCP stream. Verified against 19/19 ctest suites, the bilingual docs build, and a manually linked mqtt_interrupt_demo.

Co-Authored-By: OpenAI Codex <[email protected]>
Prove authenticated control remains deliverable while the media Connection is paused, stale Event audio is drained without refilling playback, and the next Event is accepted. Verified against 20/20 ctest suites and the focused ASan/UBSan target.

Co-Authored-By: OpenAI Codex <[email protected]>
Keep coalesced post-ack media buffered for the receive worker, ignore delayed interrupts for obsolete Events, document the cross-module ownership decisions, and add the new network-free suites to leak CI. Verified with focused TAI/control regressions.

Co-Authored-By: OpenAI Codex <[email protected]>
Route the combined regression through a TLS MQTT echo while TAI media uses the loopback peer, then cover authentication rejection, MQTT-before-media ordering, delayed TCP/MQTT duplicates, reconnect, and shutdown while pressured. Because the suite spawns mock subprocesses, keep it out of the Valgrind no-subprocess list. Verified against the focused TAI and combined regressions.

Co-Authored-By: OpenAI Codex <[email protected]>
…ordering

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, and the pinned Frame is slid only after the
remainder drains. The slide boundary is verified against the Frame's own
length field, so a metadata/storage mismatch fails fast instead of
desyncing the stream.

The MQTT control demo now handles early interrupts (before an Event's
TCP START, fail-closed unscoped interrupts, and stale Event ENDs without
terminating the Session. iot_ai_ctrl no longer rewrites a scoped
interrupt to an unscoped empty payload when payload serialization fails.

The combined regression is fully synchronized (mutex-guarded assertions),
covers same-event MQTT-before-START ordering, and records interruption
ordering inside the callback instead of inferring it post-hoc.

Verified: 20/20 tests pass serially and under ASan/UBSan (7338 assertions
in tai_integration_tests, including byte-exact multi-Packet backpressure
with per-codec-frame admission).

Co-Authored-By: MightyCode <[email protected]>
EOF
)
Use only the public TAI interface in the MQTT interruption demo and drain paused audio through its mutex-protected stale Event filter, avoiding an unsafe cross-thread SDK state mutation. Add deterministic coverage for paused, racing, queued, unscoped, and fresh Event orderings.\n\nVerified: 20/20 ctest suites pass serially; mqtt_interrupt_demo builds successfully.
Interruption filtering compared event IDs, which cannot order a stream
against an interruption that arrives out of band. Both notice paths already
carry a server time -- the MQTT asrInterrupt in its own payload, the TCP
ChatBreak in the UserData attribute (attr 111) -- and the audio media header
already carries the same clock, so the two can be compared directly.
tai_event_msg_t now exposes the borrowed user_data attribute so an
application can read it without the SDK parsing JSON.

Measured on real traffic (audio_chat_demo, real speech barge-in) the
ChatBreak cut-off lands strictly between the interrupted stream's START and
the next stream's START: 1790088542969 <= 1790088545108 < 1790088557463, all
one epoch-millisecond timeline that matched the device wall clock to the
second. The demo now discards on that comparison and fails closed onto the
in-flight stream when a notice carries no usable time, instead of ignoring
the interruption.

The paused-Packet cursor drops the fields that only duplicated state the
pause cannot change: no metadata copies, no application-side discard call,
and the storage-origin bookkeeping collapses to the pinned Frame's wire
length. A wire length that cannot be consumed exactly once fails fast with
TAI_PROTO_ERR_FRAME_DECODE rather than sliding the receive buffer out of
range (Bus error without the guard, pinned by
pending_audio_bad_wire_len).

Share one admission query across all receive checkpoints and document
when handshake media can invoke the hook. Clarify AI-control payload
re-serialization and raw-callback fallback, export its setter with IOT_API,
and correct the development SDK version. Document pacing knobs, post-ack
media delivery, and the demo's shared-test fixture contract; remove unused
loopback extern declarations and complete the Unreleased entries for #42.

Verified: full serial ctest 20/20; ASan/UBSan tai_integration_tests,
iot_tai_control_test, iot_ai_ctrl_test 3/3; docs-site i18n parity 30/30.

Co-Authored-By: MightyCode <[email protected]>
Co-Authored-By: Claude Code <[email protected]>
Co-Authored-By: OpenAI Codex <[email protected]>
@heshaoqiong-tuya
heshaoqiong-tuya force-pushed the feature/tai-backpressure2 branch from 057ffe6 to 65ce503 Compare September 29, 2026 05:45
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant