| 593 | ) |
| 594 | |
| 595 | async def recv(self): |
| 596 | # Ensure MediaClock is started |
| 597 | if not self._streamer._media_clock.is_started: |
| 598 | self._streamer._media_clock.start() |
| 599 | |
| 600 | if self.use_microphone: |
| 601 | frame = await self.audio_track.recv() |
| 602 | |
| 603 | # Apply user callback if registered |
| 604 | if self.callback is not None: |
| 605 | try: |
| 606 | processed_frame = self.callback(frame) |
| 607 | return processed_frame |
| 608 | except Exception as e: |
| 609 | self._streamer._log(f"[AUDIO] Callback error: {e}", force=True) |
| 610 | return frame |
| 611 | else: |
| 612 | return frame |
| 613 | else: |
| 614 | # Generate synthetic audio frame using MediaClock for sync |
| 615 | clock = self._streamer._media_clock |
| 616 | |
| 617 | # Create blank audio frame (silence) |
| 618 | layout = 'stereo' if self.stereo else 'mono' |
| 619 | frame = AudioFrame(format='s16', layout=layout, samples=self.samples_per_frame) |
| 620 | frame.sample_rate = self.sample_rate |
| 621 | frame.pts = clock.get_audio_pts(self._sample_count) |
| 622 | frame.time_base = fractions.Fraction(1, self.sample_rate) |
| 623 | |
| 624 | # Initialize with silence |
| 625 | bytes_per_sample = 2 # s16 = 2 bytes per sample |
| 626 | channels = 2 if self.stereo else 1 |
| 627 | for p in frame.planes: |
| 628 | p.update(bytes(self.samples_per_frame * bytes_per_sample * channels)) |
| 629 | |
| 630 | # Log first few frames |
| 631 | if self._sample_count < self.samples_per_frame * 3: |
| 632 | frame_num = self._sample_count // self.samples_per_frame |
| 633 | self._streamer._log(f"[AUDIO] Generated frame #{frame_num}: {self.samples_per_frame} samples, {frame.sample_rate}Hz, {layout}") |
| 634 | |
| 635 | # Use MediaClock for precise timing |
| 636 | sleep_time = clock.wait_until_audio_time( |
| 637 | self._sample_count + self.samples_per_frame, |
| 638 | self.sample_rate |
| 639 | ) |
| 640 | |
| 641 | # Adaptive timing: sleep if ahead, track drift if behind |
| 642 | if sleep_time > 0: |
| 643 | await asyncio.sleep(sleep_time) |
| 644 | elif sleep_time < -0.05: # More than 50ms behind |
| 645 | self._drift_corrections += 1 |
| 646 | if self._drift_corrections <= 3: |
| 647 | self._streamer._log(f"[AUDIO] Drift warning: {-sleep_time*1000:.0f}ms behind (correction #{self._drift_corrections})") |
| 648 | |
| 649 | # Apply callback if registered (callback should generate audio) |
| 650 | if self.callback is not None: |
| 651 | try: |
| 652 | processed_frame = self.callback(frame) |