Async version of poll_for_data
(self, thread_index=0, timeout=None, callback=None)
| 955 | time.sleep(0.000001) # 1μs polling interval |
| 956 | |
| 957 | async def poll_for_data_async(self, thread_index=0, timeout=None, callback=None): |
| 958 | """Async version of poll_for_data""" |
| 959 | |
| 960 | start_time = asyncio.get_event_loop().time() |
| 961 | |
| 962 | while True: |
| 963 | metadata = self._get_thread_metadata(thread_index) |
| 964 | local_tail = self.local_tails[thread_index] |
| 965 | |
| 966 | # Check if new data is available using local tail |
| 967 | if metadata["head"] != local_tail: |
| 968 | # Calculate how many new entries |
| 969 | if metadata["head"] > local_tail: |
| 970 | new_entries = metadata["head"] - local_tail |
| 971 | else: |
| 972 | new_entries = ( |
| 973 | self.config["ring_size"] - local_tail + metadata["head"] |
| 974 | ) |
| 975 | |
| 976 | data = self.consume_data(thread_index, new_entries) |
| 977 | |
| 978 | if callback: |
| 979 | for entry_data in data: |
| 980 | await callback(entry_data) |
| 981 | else: |
| 982 | return data |
| 983 | |
| 984 | if timeout and (asyncio.get_event_loop().time() - start_time) > timeout: |
| 985 | return [] |
| 986 | |
| 987 | # Yield control to asyncio |
| 988 | await asyncio.sleep(0.000001) # 1μs polling interval |
| 989 | |
| 990 | def get_config(self): |
| 991 | """Get ring buffer configuration""" |
no test coverage detected