| 27 | |
| 28 | |
| 29 | class ThreadedRecordWriter: |
| 30 | def __init__(self, writer): |
| 31 | self._read_q = queue.Queue() |
| 32 | self._thread = threading.Thread( |
| 33 | target=self._threaded_record_writer, args=(writer,) |
| 34 | ) |
| 35 | |
| 36 | def _threaded_record_writer(self, writer): |
| 37 | while True: |
| 38 | record = self._read_q.get() |
| 39 | if record is False: |
| 40 | return |
| 41 | writer.write_record(record) |
| 42 | |
| 43 | def write_record(self, record): |
| 44 | self._read_q.put_nowait(record) |
| 45 | |
| 46 | def start(self): |
| 47 | self._thread.start() |
| 48 | |
| 49 | def close(self): |
| 50 | self._read_q.put_nowait(False) |
| 51 | self._thread.join() |
| 52 | |
| 53 | |
| 54 | class BaseDatabaseTest(unittest.TestCase): |
no outgoing calls