refactored
This commit is contained in:
+20
-21
@@ -1,58 +1,57 @@
|
||||
import socket
|
||||
import selectors
|
||||
from copy import deepcopy
|
||||
import time
|
||||
|
||||
|
||||
class MultiClient(object):
|
||||
def __init__(self, servers: list[dict]):
|
||||
self.servers = servers
|
||||
self.sock_list = deepcopy(self.servers)
|
||||
self.sock_list = [{'name': s['name'], 'addr': (s['host'], s['port'])} for s in servers]
|
||||
self.sel = None
|
||||
|
||||
for server in self.sock_list:
|
||||
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
sock.setblocking(False)
|
||||
server['sock'] = sock
|
||||
|
||||
def find_by_addr(self, addr):
|
||||
for sock in self.sock_list:
|
||||
if sock['addr'] == addr:
|
||||
return sock
|
||||
return None
|
||||
|
||||
def connect(self):
|
||||
self.sel = selectors.DefaultSelector()
|
||||
events = selectors.EVENT_READ
|
||||
for sock in self.sock_list:
|
||||
print(f"Starting connection {sock['name']} to {sock['host']}:{sock['port']}")
|
||||
print(f"Starting connection {sock['name']} to {sock['addr']}")
|
||||
self.sel.register(sock['sock'], events, data=None)
|
||||
sock['sock'].connect_ex((sock['host'], sock['port']))
|
||||
sock['sock'].connect_ex(sock['addr'])
|
||||
|
||||
def disconnect(self):
|
||||
for sock in self.sock_list:
|
||||
sock['sock'].close()
|
||||
|
||||
def on_receive(self, key):
|
||||
data = key.fileobj.recv(1024)
|
||||
print(f"{self}: on_receive: {data}")
|
||||
|
||||
def on_send(self, key):
|
||||
return None
|
||||
def on_receive(self, rx_list: list):
|
||||
for data in rx_list:
|
||||
print(f"{self}: on_receive: 'Time': {data['timestamp'], data['name']}")
|
||||
|
||||
def send_request(self):
|
||||
events = selectors.EVENT_READ | selectors.EVENT_WRITE
|
||||
for sock in self.sock_list:
|
||||
self.sel.register(sock['sock'], events, data=kakac)
|
||||
self.sel.register(sock['sock'], events, data=None)
|
||||
print("receive")
|
||||
|
||||
def event_loop(self):
|
||||
try:
|
||||
while True:
|
||||
events = self.sel.select(timeout=1)
|
||||
rx_data = []
|
||||
for key, mask in events:
|
||||
descr = self.find_by_addr(key.fileobj.getpeername())
|
||||
if mask & selectors.EVENT_READ:
|
||||
self.on_receive(key)
|
||||
if mask & selectors.EVENT_WRITE:
|
||||
data = self.on_send(key)
|
||||
if data is None:
|
||||
events = selectors.EVENT_READ
|
||||
for sock in self.sock_list:
|
||||
self.sel.register(sock['sock'], events, data=None)
|
||||
timestamp = time.time()
|
||||
data = key.fileobj.recv(1024)
|
||||
rx_data.append({'timestamp': timestamp, 'name': descr['name'], 'data': data})
|
||||
|
||||
self.on_receive(rx_data)
|
||||
|
||||
except KeyboardInterrupt:
|
||||
print("Caught keyboard interrupt, exiting")
|
||||
|
||||
Reference in New Issue
Block a user