audio_receiver.c 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644
  1. #include <inttypes.h>
  2. #include <stdlib.h>
  3. #include <string.h>
  4. #include "audio_receiver.h"
  5. #include "esp_heap_caps.h"
  6. #include "esp_log.h"
  7. #include "esp_timer.h"
  8. #include "audio_buffer.h"
  9. #include "audio_decoder.h"
  10. #include "audio_output.h"
  11. #include "audio_receiver_internal.h"
  12. #include "audio_stream.h"
  13. #include "audio_timing.h"
  14. #include "ptp_clock.h"
  15. #define DEFAULT_SAMPLE_RATE 44100
  16. #define DEFAULT_CHANNELS 2
  17. #define DEFAULT_BITS_PER_SAMPLE 16
  18. #define DEFAULT_FRAME_SIZE 352
  19. #define DECRYPT_BUFFER_SIZE 8192
  20. static const char *TAG = "audio_recv";
  21. static audio_receiver_state_t receiver = {0};
  22. static void audio_receiver_reset_stats(void) {
  23. memset(&receiver.stats, 0, sizeof(receiver.stats));
  24. }
  25. static void audio_receiver_reset_resend_state(void) {
  26. receiver.rtp_sequence_valid = false;
  27. receiver.resend_window_first = 0;
  28. receiver.resend_missing_mask = 0;
  29. receiver.resend_last_request_time_us = 0;
  30. receiver.last_resend_error_time_us = 0;
  31. }
  32. static void audio_receiver_reset_blocks(void) {
  33. receiver.blocks_read = 0;
  34. receiver.blocks_read_in_sequence = 0;
  35. }
  36. static void audio_receiver_copy_stream_state(audio_stream_t *dst,
  37. const audio_stream_t *src) {
  38. if (!dst || !src) {
  39. return;
  40. }
  41. dst->format = src->format;
  42. dst->encrypt = src->encrypt;
  43. }
  44. esp_err_t audio_receiver_init(void) {
  45. if (receiver.buffer.pool) {
  46. return ESP_OK;
  47. }
  48. receiver.realtime_stream = audio_stream_create_realtime();
  49. if (receiver.realtime_stream) {
  50. receiver.realtime_stream->ctx = &receiver;
  51. }
  52. receiver.buffered_stream = audio_stream_create_buffered();
  53. if (receiver.buffered_stream) {
  54. receiver.buffered_stream->ctx = &receiver;
  55. }
  56. if (!receiver.realtime_stream || !receiver.buffered_stream) {
  57. ESP_LOGE(TAG, "Failed to allocate audio streams");
  58. audio_stream_destroy(receiver.realtime_stream);
  59. audio_stream_destroy(receiver.buffered_stream);
  60. receiver.realtime_stream = NULL;
  61. receiver.buffered_stream = NULL;
  62. return ESP_ERR_NO_MEM;
  63. }
  64. receiver.stream = receiver.realtime_stream;
  65. audio_format_t default_format = {0};
  66. strcpy(default_format.codec, "AppleLossless");
  67. default_format.sample_rate = DEFAULT_SAMPLE_RATE;
  68. default_format.channels = DEFAULT_CHANNELS;
  69. default_format.bits_per_sample = DEFAULT_BITS_PER_SAMPLE;
  70. default_format.frame_size = DEFAULT_FRAME_SIZE;
  71. receiver.realtime_stream->format = default_format;
  72. receiver.buffered_stream->format = default_format;
  73. esp_err_t err = audio_buffer_init(&receiver.buffer);
  74. if (err != ESP_OK) {
  75. audio_stream_destroy(receiver.realtime_stream);
  76. audio_stream_destroy(receiver.buffered_stream);
  77. receiver.realtime_stream = NULL;
  78. receiver.buffered_stream = NULL;
  79. return err;
  80. }
  81. receiver.decrypt_buffer_size = DECRYPT_BUFFER_SIZE;
  82. #ifdef CONFIG_SPIRAM
  83. receiver.decrypt_buffer = heap_caps_malloc(
  84. receiver.decrypt_buffer_size, MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT);
  85. #endif
  86. if (!receiver.decrypt_buffer) {
  87. receiver.decrypt_buffer = malloc(receiver.decrypt_buffer_size);
  88. }
  89. if (!receiver.decrypt_buffer) {
  90. ESP_LOGE(TAG, "Failed to allocate decrypt buffer");
  91. audio_buffer_deinit(&receiver.buffer);
  92. audio_stream_destroy(receiver.realtime_stream);
  93. audio_stream_destroy(receiver.buffered_stream);
  94. receiver.realtime_stream = NULL;
  95. receiver.buffered_stream = NULL;
  96. return ESP_ERR_NO_MEM;
  97. }
  98. size_t pending_capacity =
  99. sizeof(audio_frame_header_t) +
  100. ((size_t)MAX_SAMPLES_PER_FRAME * AUDIO_MAX_CHANNELS * sizeof(int16_t));
  101. audio_timing_init(&receiver.timing, pending_capacity);
  102. audio_timing_set_format(&receiver.timing, &receiver.stream->format);
  103. receiver.buffered_listen_socket = -1;
  104. receiver.buffered_client_socket = -1;
  105. audio_receiver_reset_blocks();
  106. return ESP_OK;
  107. }
  108. void audio_receiver_set_format(const audio_format_t *format) {
  109. if (!format) {
  110. return;
  111. }
  112. if (!receiver.realtime_stream || !receiver.buffered_stream) {
  113. return;
  114. }
  115. receiver.realtime_stream->format = *format;
  116. receiver.buffered_stream->format = *format;
  117. audio_decoder_destroy(receiver.decoder);
  118. receiver.decoder = NULL;
  119. audio_decoder_config_t cfg = {.format = *format};
  120. receiver.decoder = audio_decoder_create(&cfg);
  121. if (!receiver.decoder) {
  122. ESP_LOGW(TAG, "Decoder not initialized for codec: %s", format->codec);
  123. }
  124. audio_timing_set_format(&receiver.timing, format);
  125. audio_output_set_source_rate(format->sample_rate);
  126. }
  127. void audio_receiver_set_encryption(const audio_encrypt_t *encrypt) {
  128. if (!receiver.realtime_stream || !receiver.buffered_stream) {
  129. return;
  130. }
  131. if (encrypt) {
  132. receiver.realtime_stream->encrypt = *encrypt;
  133. receiver.buffered_stream->encrypt = *encrypt;
  134. } else {
  135. memset(&receiver.realtime_stream->encrypt, 0,
  136. sizeof(receiver.realtime_stream->encrypt));
  137. memset(&receiver.buffered_stream->encrypt, 0,
  138. sizeof(receiver.buffered_stream->encrypt));
  139. }
  140. }
  141. void audio_receiver_set_output_latency_us(uint32_t latency_us) {
  142. if (!receiver.stream) {
  143. return;
  144. }
  145. audio_timing_set_output_latency(&receiver.timing, &receiver.stream->format,
  146. latency_us);
  147. }
  148. uint32_t audio_receiver_get_output_latency_us(void) {
  149. return audio_timing_get_output_latency(&receiver.timing);
  150. }
  151. uint32_t audio_receiver_get_hardware_latency_us(void) {
  152. return audio_timing_get_hardware_latency();
  153. }
  154. uint32_t audio_receiver_get_advertised_latency_us(void) {
  155. return audio_timing_get_advertised_latency(&receiver.timing);
  156. }
  157. void audio_receiver_set_anchor_time(uint64_t clock_id, uint64_t network_time_ns,
  158. uint32_t rtp_time) {
  159. if (!receiver.stream) {
  160. return;
  161. }
  162. int sample_rate = receiver.stream->format.sample_rate;
  163. if (sample_rate <= 0) {
  164. sample_rate = 44100;
  165. }
  166. // Window size for the upper RTP gate: 10 s of samples. Large enough that
  167. // a normal 2-4 s pre-buffer passes, but small enough to reject stale frames
  168. // left in the TCP socket buffer after a backward seek.
  169. const uint32_t gate_window = (uint32_t)(10 * sample_rate);
  170. const int32_t seek_threshold = 5 * sample_rate;
  171. // --- Phase 1: Arm RTP gates BEFORE opening the blanket gate -----------
  172. //
  173. // The blanket gate (discard_all_until_anchor) blocks ALL frames from the
  174. // TCP task. The per-RTP gates filter by timestamp range. On single-core
  175. // ESP32-S2, ESP_LOGI can yield to the scheduler, so any gap between
  176. // clearing the blanket and arming the per-RTP gates lets the TCP task
  177. // queue stale frames. Arm first, then open.
  178. bool gates_armed = false;
  179. // Path A: seek_flush set arm_gate_on_next_anchor because the buffer was
  180. // already empty when the flush happened (forward-seek).
  181. if (receiver.arm_gate_on_next_anchor) {
  182. receiver.arm_gate_on_next_anchor = false;
  183. receiver.discard_before_rtp = rtp_time;
  184. receiver.discard_before_rtp_valid = true;
  185. receiver.discard_above_rtp = rtp_time + gate_window;
  186. receiver.discard_above_rtp_valid = true;
  187. gates_armed = true;
  188. ESP_LOGI(TAG,
  189. "RTP gates armed on anchor: discard_before=%lu discard_above=%lu",
  190. (unsigned long)rtp_time, (unsigned long)(rtp_time + gate_window));
  191. }
  192. // Path B: Anchor-change detection — the phone changed track with a
  193. // PAUSE → RESUME cycle but no FLUSHBUFFERED. The buffer may already be
  194. // empty (consumed during playback), so the seek-detection heuristic
  195. // below (which needs oldest_rtp from the buffer) would miss it.
  196. //
  197. // Compare the new anchor against the EXPECTED current playback position
  198. // (old_anchor_rtp + elapsed_time × sample_rate), NOT against the raw old
  199. // anchor. The raw old anchor was set at the start of the previous play
  200. // segment; comparing against it gives a delta equal to (elapsed_play_time +
  201. // new_anchor_lead_time), which easily exceeds the 5-second threshold on a
  202. // normal pause/resume within the same track — causing a false flush that
  203. // empties valid pre-buffered frames and produces 6+ seconds of silence.
  204. // Using the expected position instead, normal resume gives a delta of only
  205. // the anchor's lead-time offset (< 2 s), while a real track-change or seek
  206. // gives a huge delta (many minutes).
  207. if (!gates_armed && receiver.timing.anchor_valid) {
  208. // Choose reference point for the seek-detection comparison:
  209. // - If we have a pause snapshot, use it. The snapshot was taken at the
  210. // exact moment the sender said PAUSE, so it reflects the true pause
  211. // position rather than a wall-clock estimate that keeps running during
  212. // the pause and overshoots by (pause_duration x sample_rate).
  213. // - Otherwise fall back to the elapsed-time estimate (covers the edge
  214. // case where a track changes without a prior PAUSE signal).
  215. uint32_t reference_rtp;
  216. if (receiver.paused_rtp_valid) {
  217. reference_rtp = receiver.paused_rtp;
  218. receiver.paused_rtp_valid = false; // one-shot: consume after use
  219. // Compute the pause duration from the RTP snapshot so we can notify
  220. // the PTP clock without tracking a separate wall-clock timestamp.
  221. // anchor_local_time_ns/1000 is the µs when the anchor was set;
  222. // adding the played-sample offset gives the µs when play paused.
  223. int32_t played =
  224. (int32_t)(reference_rtp - receiver.timing.anchor_rtp_time);
  225. int64_t pause_time_us = receiver.timing.anchor_local_time_ns / 1000LL +
  226. (int64_t)played * 1000000LL / sample_rate;
  227. int64_t pause_us = esp_timer_get_time() - pause_time_us;
  228. ptp_clock_notify_resume((pause_us > 0) ? (uint32_t)(pause_us / 1000LL)
  229. : 0);
  230. ESP_LOGD(TAG, "Path B: pause snapshot rtp=%lu pause=%.1f s",
  231. (unsigned long)reference_rtp, (float)pause_us / 1e6f);
  232. } else {
  233. int64_t elapsed_us = esp_timer_get_time() -
  234. (receiver.timing.anchor_local_time_ns / 1000LL);
  235. if (elapsed_us < 0) {
  236. elapsed_us = 0;
  237. }
  238. // Cap elapsed to prevent int64 overflow on very long pauses.
  239. if (elapsed_us > 600000000LL) {
  240. elapsed_us = 600000000LL;
  241. }
  242. int32_t elapsed_samples =
  243. (int32_t)((elapsed_us * (int64_t)sample_rate) / 1000000LL);
  244. reference_rtp =
  245. receiver.timing.anchor_rtp_time + (uint32_t)elapsed_samples;
  246. }
  247. int32_t delta = (int32_t)(rtp_time - reference_rtp);
  248. int32_t abs_delta = delta < 0 ? -delta : delta;
  249. if (abs_delta > seek_threshold) {
  250. ESP_LOGI(TAG,
  251. "Anchor change detected: ref_rtp=%lu new_rtp=%lu "
  252. "delta=%ld samples (%.1f s) - flushing & arming gates",
  253. (unsigned long)reference_rtp, (unsigned long)rtp_time,
  254. (long)delta, (float)delta / sample_rate);
  255. audio_buffer_flush(&receiver.buffer);
  256. receiver.timing.playout_started = false;
  257. receiver.timing.pending_valid = false;
  258. receiver.timing.pending_frame_len = 0;
  259. receiver.timing.ready_time_us = 0;
  260. receiver.timing.deferred_flush_pending = false;
  261. receiver.blocks_read_in_sequence = 0;
  262. receiver.discard_before_rtp = rtp_time;
  263. receiver.discard_before_rtp_valid = true;
  264. receiver.discard_above_rtp = rtp_time + gate_window;
  265. receiver.discard_above_rtp_valid = true;
  266. receiver.timing.quick_start = true;
  267. gates_armed = true;
  268. } else {
  269. ESP_LOGD(TAG,
  270. "Anchor resume OK: ref_rtp=%lu new_rtp=%lu "
  271. "delta=%ld samples (%.2f s) - same track, no flush",
  272. (unsigned long)reference_rtp, (unsigned long)rtp_time,
  273. (long)delta, (float)delta / (float)sample_rate);
  274. }
  275. }
  276. // NOW safe to clear the blanket gate — per-RTP gates are active.
  277. receiver.discard_all_until_anchor = false;
  278. // --- Phase 2: Seek detection from buffer content ----------------------
  279. //
  280. // If stale data managed to enter the buffer (e.g. queued before
  281. // seek_flush was called), detect it by comparing the oldest buffered
  282. // RTP timestamp against the new anchor.
  283. uint32_t oldest_rtp = 0;
  284. if (audio_buffer_oldest_timestamp(&receiver.buffer, &oldest_rtp)) {
  285. int32_t rtp_ahead = (int32_t)(oldest_rtp - rtp_time);
  286. int32_t abs_ahead = rtp_ahead < 0 ? -rtp_ahead : rtp_ahead;
  287. if (abs_ahead > seek_threshold) {
  288. ESP_LOGI(TAG,
  289. "Seek detected: oldest_rtp=%lu, new anchor rtp=%lu, "
  290. "delta=%ld samples (%.1f s) — flushing stale buffer",
  291. (unsigned long)oldest_rtp, (unsigned long)rtp_time,
  292. (long)rtp_ahead, (float)rtp_ahead / sample_rate);
  293. audio_buffer_flush(&receiver.buffer);
  294. receiver.timing.playout_started = false;
  295. receiver.timing.pending_valid = false;
  296. receiver.timing.pending_frame_len = 0;
  297. receiver.timing.ready_time_us = 0;
  298. receiver.timing.deferred_flush_pending = false;
  299. receiver.blocks_read_in_sequence = 0;
  300. receiver.timing.quick_start = true;
  301. if (!gates_armed) {
  302. receiver.discard_before_rtp = rtp_time;
  303. receiver.discard_before_rtp_valid = true;
  304. receiver.discard_above_rtp = rtp_time + gate_window;
  305. receiver.discard_above_rtp_valid = true;
  306. }
  307. }
  308. }
  309. // Pin the PTP clock to the master announced by the anchor packet's
  310. // clock_id field. Without this, ptp_clock can lock to any PTP master
  311. // on the LAN (HomePods, AppleTVs, NTP-PTP gateways) and produce offsets
  312. // that have nothing to do with the AirPlay sender's clock domain.
  313. // clock_id == 0 happens on the AirPlay 1 NTP path; in that case leave
  314. // the filter as-is (set or cleared by a prior 0xD7 anchor / TEARDOWN).
  315. if (clock_id != 0) {
  316. ptp_clock_set_master_clock_id(clock_id);
  317. }
  318. audio_timing_set_anchor(&receiver.timing, &receiver.stream->format, clock_id,
  319. network_time_ns, rtp_time);
  320. }
  321. void audio_receiver_set_playing(bool playing) {
  322. audio_timing_set_playing(&receiver.timing, playing);
  323. if (!playing) {
  324. receiver.blocks_read_in_sequence = 0;
  325. // Snapshot the expected RTP position at the moment of pause so that
  326. // Path B in audio_receiver_set_anchor_time() can compare the next
  327. // resume anchor against the actual pause position.
  328. //
  329. // Without this, Path B uses (anchor_rtp + wall_clock_elapsed), which
  330. // overshoots by the pause duration and fires a false seek flush on any
  331. // pause >= seek_threshold (5 s) — causing up to 7+ s of silence when
  332. // pre-buffered frames end up far ahead of the unwanted new anchor.
  333. if (receiver.timing.anchor_valid && receiver.stream) {
  334. int sample_rate = receiver.stream->format.sample_rate;
  335. if (sample_rate <= 0) {
  336. sample_rate = 44100;
  337. }
  338. int64_t elapsed_us = esp_timer_get_time() -
  339. (receiver.timing.anchor_local_time_ns / 1000LL);
  340. if (elapsed_us < 0) {
  341. elapsed_us = 0;
  342. }
  343. if (elapsed_us > 600000000LL) {
  344. elapsed_us = 600000000LL;
  345. }
  346. int32_t elapsed_samples =
  347. (int32_t)((elapsed_us * (int64_t)sample_rate) / 1000000LL);
  348. receiver.paused_rtp =
  349. receiver.timing.anchor_rtp_time + (uint32_t)elapsed_samples;
  350. receiver.paused_rtp_valid = true;
  351. ESP_LOGD(TAG, "Pause: RTP snapshot=%lu (elapsed=%.2f s)",
  352. (unsigned long)receiver.paused_rtp, (float)elapsed_us / 1e6f);
  353. }
  354. }
  355. }
  356. void audio_receiver_reset_timing(void) {
  357. audio_timing_reset(&receiver.timing);
  358. }
  359. bool audio_receiver_is_playing(void) {
  360. return receiver.timing.playing;
  361. }
  362. void audio_receiver_set_stream_type(audio_stream_type_t type) {
  363. if (!receiver.realtime_stream || !receiver.buffered_stream) {
  364. return;
  365. }
  366. audio_stream_t *target = receiver.realtime_stream;
  367. if (type == AUDIO_STREAM_BUFFERED) {
  368. target = receiver.buffered_stream;
  369. }
  370. if (!target) {
  371. return;
  372. }
  373. if (receiver.stream != target) {
  374. if (receiver.stream) {
  375. audio_receiver_copy_stream_state(target, receiver.stream);
  376. if (receiver.stream->running && receiver.stream->ops &&
  377. receiver.stream->ops->stop) {
  378. receiver.stream->ops->stop(receiver.stream);
  379. }
  380. }
  381. receiver.stream = target;
  382. }
  383. receiver.stream->type = type;
  384. }
  385. esp_err_t audio_receiver_start(uint16_t data_port, uint16_t control_port) {
  386. audio_receiver_set_stream_type(AUDIO_STREAM_REALTIME);
  387. if (!receiver.stream || !receiver.stream->ops ||
  388. !receiver.stream->ops->start) {
  389. return ESP_FAIL;
  390. }
  391. // Always stop and restart fresh
  392. if (receiver.stream->running) {
  393. receiver.stream->ops->stop(receiver.stream);
  394. }
  395. receiver.data_port = data_port;
  396. receiver.control_port = control_port;
  397. // Starting a stream resets all timing state (including pause tracking)
  398. audio_receiver_reset_stats();
  399. audio_buffer_flush(&receiver.buffer);
  400. audio_timing_reset(&receiver.timing);
  401. audio_receiver_reset_resend_state();
  402. receiver.timing.ptp_locked = ptp_clock_is_locked();
  403. audio_receiver_reset_blocks();
  404. return receiver.stream->ops->start(receiver.stream, data_port);
  405. }
  406. esp_err_t audio_receiver_start_buffered(uint16_t tcp_port) {
  407. audio_receiver_set_stream_type(AUDIO_STREAM_BUFFERED);
  408. if (!receiver.stream || !receiver.stream->ops ||
  409. !receiver.stream->ops->start) {
  410. return ESP_FAIL;
  411. }
  412. // Buffered streams use a fixed port, no need to restart if running
  413. if (receiver.stream->running) {
  414. return ESP_OK;
  415. }
  416. // Starting a stream resets all timing state (including pause tracking)
  417. audio_receiver_reset_stats();
  418. audio_buffer_flush(&receiver.buffer);
  419. audio_timing_reset(&receiver.timing);
  420. audio_receiver_reset_resend_state();
  421. receiver.timing.ptp_locked = ptp_clock_is_locked();
  422. audio_receiver_reset_blocks();
  423. return receiver.stream->ops->start(receiver.stream, tcp_port);
  424. }
  425. esp_err_t audio_receiver_start_stream(uint16_t data_port, uint16_t control_port,
  426. uint16_t tcp_port) {
  427. if (!receiver.stream) {
  428. return ESP_FAIL;
  429. }
  430. if (receiver.stream->type == AUDIO_STREAM_BUFFERED) {
  431. return audio_receiver_start_buffered(tcp_port);
  432. }
  433. return audio_receiver_start(data_port, control_port);
  434. }
  435. uint16_t audio_receiver_get_stream_port(void) {
  436. if (!receiver.stream || !receiver.stream->ops ||
  437. !receiver.stream->ops->get_port) {
  438. return 0;
  439. }
  440. return receiver.stream->ops->get_port(receiver.stream);
  441. }
  442. void audio_receiver_set_client_control(uint32_t client_ip,
  443. uint16_t client_control_port) {
  444. if (client_ip == 0 || client_control_port == 0) {
  445. receiver.retransmit_enabled = false;
  446. audio_receiver_reset_resend_state();
  447. return;
  448. }
  449. memset(&receiver.client_control_addr, 0,
  450. sizeof(receiver.client_control_addr));
  451. receiver.client_control_addr.sin_family = AF_INET;
  452. receiver.client_control_addr.sin_addr.s_addr = client_ip;
  453. receiver.client_control_addr.sin_port = htons(client_control_port);
  454. receiver.retransmit_enabled = true;
  455. audio_receiver_reset_resend_state();
  456. ESP_LOGI(TAG, "NACK retransmission enabled, client control port %u",
  457. client_control_port);
  458. }
  459. void audio_receiver_stop(void) {
  460. if (receiver.realtime_stream && receiver.realtime_stream->ops &&
  461. receiver.realtime_stream->ops->stop) {
  462. receiver.realtime_stream->ops->stop(receiver.realtime_stream);
  463. }
  464. if (receiver.buffered_stream && receiver.buffered_stream->ops &&
  465. receiver.buffered_stream->ops->stop) {
  466. receiver.buffered_stream->ops->stop(receiver.buffered_stream);
  467. }
  468. audio_decoder_destroy(receiver.decoder);
  469. receiver.decoder = NULL;
  470. if (receiver.realtime_stream) {
  471. memset(&receiver.realtime_stream->encrypt, 0,
  472. sizeof(receiver.realtime_stream->encrypt));
  473. }
  474. if (receiver.buffered_stream) {
  475. memset(&receiver.buffered_stream->encrypt, 0,
  476. sizeof(receiver.buffered_stream->encrypt));
  477. }
  478. receiver.retransmit_enabled = false;
  479. memset(&receiver.client_control_addr, 0,
  480. sizeof(receiver.client_control_addr));
  481. audio_receiver_reset_resend_state();
  482. audio_receiver_flush();
  483. }
  484. void audio_receiver_stop_buffered_only(void) {
  485. if (receiver.buffered_stream && receiver.buffered_stream->ops &&
  486. receiver.buffered_stream->ops->stop) {
  487. receiver.buffered_stream->ops->stop(receiver.buffered_stream);
  488. }
  489. }
  490. void audio_receiver_get_stats(audio_stats_t *stats) {
  491. if (!stats) {
  492. return;
  493. }
  494. memcpy(stats, &receiver.stats, sizeof(receiver.stats));
  495. }
  496. size_t audio_receiver_read(int16_t *buffer, size_t samples) {
  497. if (!receiver.buffer.pool || !buffer || samples == 0) {
  498. return 0;
  499. }
  500. return audio_timing_read(&receiver.timing, &receiver.buffer, receiver.stream,
  501. &receiver.stats, buffer, samples);
  502. }
  503. bool audio_receiver_has_data(void) {
  504. int buffered_frames = audio_buffer_get_frame_count(&receiver.buffer);
  505. return buffered_frames > 0 || receiver.timing.pending_valid;
  506. }
  507. void audio_receiver_flush(void) {
  508. // Flush is an explicit reset — clear all timing state including pause
  509. // tracking. The sender will provide fresh anchor times after flush.
  510. // Also disarm any pending deferred flush so it does not fire on the
  511. // next track's frames.
  512. audio_buffer_flush(&receiver.buffer);
  513. audio_timing_reset(&receiver.timing);
  514. audio_receiver_reset_resend_state();
  515. receiver.discard_before_rtp_valid = false;
  516. receiver.discard_above_rtp_valid = false;
  517. receiver.arm_gate_on_next_anchor = false;
  518. receiver.discard_all_until_anchor = false;
  519. receiver.paused_rtp_valid = false;
  520. receiver.blocks_read_in_sequence = 1;
  521. }
  522. void audio_receiver_seek_flush(void) {
  523. // Mid-stream seek flush (FLUSH / immediate FLUSHBUFFERED). Like
  524. // audio_receiver_flush() but sets timing.quick_start so audio_timing_read
  525. // starts as soon as 1 frame is available, with normal anchor-based timing.
  526. // Also disarms any pending deferred flush (audio_timing_reset clears it).
  527. audio_receiver_flush();
  528. receiver.timing.quick_start = true;
  529. // Request that the RTP gate be armed as soon as the next anchor arrives.
  530. // This covers the forward-seek case where the buffer is already empty by
  531. // the time SETRATEANCHORTIME arrives, so the seek-detection heuristic
  532. // (which needs oldest_rtp from the buffer) would otherwise miss arming it.
  533. receiver.arm_gate_on_next_anchor = true;
  534. // Reject ALL incoming frames until the next anchor. Prevents stale TCP
  535. // data from filling the buffer between FLUSHBUFFERED and SETRATEANCHORTIME,
  536. // which would cause a second flush and double the startup delay.
  537. receiver.discard_all_until_anchor = true;
  538. }
  539. void audio_receiver_set_deferred_flush(uint32_t flush_until_ts) {
  540. if (!receiver.stream) {
  541. return;
  542. }
  543. // Write flush_until_ts before arming the flag so audio_timing_read never
  544. // sees deferred_flush_pending=true with a stale timestamp.
  545. receiver.timing.flush_until_ts = flush_until_ts;
  546. receiver.timing.deferred_flush_pending = true;
  547. ESP_LOGI(TAG, "Deferred flush armed: flush_until_ts=%" PRIu32,
  548. flush_until_ts);
  549. }
  550. void audio_receiver_pause(void) {
  551. // Stop the consumer. The receiver tasks keep running so the audio buffer
  552. // continues to fill with pre-buffered audio — TCP back-pressure naturally
  553. // throttles the sender. On resume the phone sends a fresh
  554. // SETRATEANCHORTIME anchor that re-aligns the buffered frames to the
  555. // correct wall-clock position; no flush or offset compensation is needed.
  556. audio_timing_set_playing(&receiver.timing, false);
  557. receiver.blocks_read_in_sequence = 0;
  558. }
  559. uint16_t audio_receiver_get_buffered_port(void) {
  560. return receiver.buffered_port;
  561. }