Initialize a new channel stream
(self, url: str, channel_id: str)
| 807 | self.user_agent: str = user_agent or Config.DEFAULT_USER_AGENT |
| 808 | |
| 809 | def initialize_channel(self, url: str, channel_id: str) -> None: |
| 810 | """Initialize a new channel stream""" |
| 811 | if channel_id in self.stream_managers: |
| 812 | self.stop_channel(channel_id) |
| 813 | |
| 814 | self.stream_managers[channel_id] = StreamManager( |
| 815 | url, |
| 816 | channel_id, |
| 817 | user_agent=self.user_agent |
| 818 | ) |
| 819 | self.stream_buffers[channel_id] = StreamBuffer() |
| 820 | self.client_managers[channel_id] = ClientManager() |
| 821 | |
| 822 | # Set up cleanup references |
| 823 | self.stream_managers[channel_id].client_manager = self.client_managers[channel_id] |
| 824 | self.stream_managers[channel_id].proxy_server = self |
| 825 | |
| 826 | fetcher = StreamFetcher( |
| 827 | self.stream_managers[channel_id], |
| 828 | self.stream_buffers[channel_id] |
| 829 | ) |
| 830 | |
| 831 | self.fetch_threads[channel_id] = threading.Thread( |
| 832 | target=fetcher.fetch_loop, |
| 833 | name=f"StreamFetcher-{channel_id}", |
| 834 | daemon=True |
| 835 | ) |
| 836 | self.fetch_threads[channel_id].start() |
| 837 | |
| 838 | # Start cleanup monitoring |
| 839 | self.stream_managers[channel_id].start_cleanup_thread() |
| 840 | logging.info(f"Initialized channel {channel_id} with URL {url}") |
| 841 | |
| 842 | def stop_channel(self, channel_id: str) -> None: |
| 843 | """Stop and cleanup a channel""" |
no test coverage detected