节点队列同步线程
(self)
| 116 | logger.info(f"[ Global GID: {self.current_gid} | SyncPool Size: {len(self.sync_pool)} | {self.total_ack} acked / {self.recv_count} total ] Current {self.send_count} sent, {self.ack_advantages} ack. Speed {(self.recv_count-last_count)/interval:.2f}/s.") |
| 117 | |
| 118 | def _sync_node_queue(self): |
| 119 | """节点队列同步线程""" |
| 120 | while True: |
| 121 | time.sleep(3) |
| 122 | with self.sync_lock: |
| 123 | self.sync_sender.send_multipart([ |
| 124 | b"SYNC_NODE_QUEUE_LENGTHS", |
| 125 | pickle.dumps(self.node_queue_lengths) |
| 126 | ]) |
| 127 | logger.debug(f"Sync node queue lengths {len(self.node_queue_lengths)}") |
| 128 | |
| 129 | def start(self): |
| 130 | """启动同步管理器""" |
nothing calls this directly
no outgoing calls
no test coverage detected