1 """ TCP server that communicates with terminals """
3 from configparser import ConfigParser
4 from importlib import import_module
5 from logging import getLogger
15 from struct import pack
17 from typing import Any, cast, Dict, List, Optional, Set, Tuple, Union
21 from .zmsg import Bcast, Resp
23 log = getLogger("loctrkd/collector")
30 def recv(self, segment: bytes) -> List[Union[bytes, str]]:
33 def close(self) -> bytes:
37 def enframe(buffer: bytes, imei: Optional[str] = None) -> bytes:
41 def probe_buffer(buffer: bytes) -> bool:
45 def parse_message(packet: bytes, is_incoming: bool = True) -> Any:
49 def inline_response(packet: bytes) -> Optional[bytes]:
53 def is_goodbye_packet(packet: bytes) -> bool:
57 def imei_from_packet(packet: bytes) -> Optional[str]:
61 def proto_of_message(packet: bytes) -> str:
65 pmods: List[ProtoModule] = []
69 """Connected socket to the terminal plus buffer and metadata"""
71 def __init__(self, sock: socket, addr: Any) -> None:
74 self.pmod: Optional[ProtoModule] = None
75 self.stream: Optional[ProtoModule.Stream] = None
76 self.imei: Optional[str] = None
78 def close(self) -> None:
79 log.debug("Closing fd %d (IMEI %s)", self.sock.fileno(), self.imei)
82 rest = self.stream.close()
87 "%d bytes in buffer on close: %s", len(rest), rest[:64].hex()
90 def recv(self) -> Optional[List[Tuple[float, Any, bytes]]]:
91 """Read from the socket and parse complete messages"""
93 segment = self.sock.recv(MAXBUFFER)
96 "Reading from fd %d (IMEI %s): %s",
102 if not segment: # Terminal has closed connection
104 "EOF reading from fd %d (IMEI %s)",
109 if self.stream is None:
111 if pmod.probe_buffer(segment):
113 self.stream = pmod.Stream()
115 if self.stream is None:
117 "unrecognizable %d bytes of data %s from fd %d",
125 for elem in self.stream.recv(segment):
126 if isinstance(elem, bytes):
127 msgs.append((when, self.addr, elem))
130 "%s from fd %d (IMEI %s)",
137 def send(self, buffer: bytes) -> None:
138 assert self.stream is not None and self.pmod is not None
140 self.sock.send(self.pmod.enframe(buffer, imei=self.imei))
143 "Sending to fd %d (IMEI %s): %s",
151 def __init__(self) -> None:
152 self.by_fd: Dict[int, Client] = {}
153 self.by_imei: Dict[str, Client] = {}
155 def fds(self) -> Set[int]:
156 return set(self.by_fd.keys())
158 def add(self, clntsock: socket, clntaddr: Any) -> int:
159 fd = clntsock.fileno()
160 log.info("Start serving fd %d from %s", fd, clntaddr)
161 self.by_fd[fd] = Client(clntsock, clntaddr)
164 def stop(self, fd: int) -> None:
165 if fd not in self.by_fd:
166 log.debug("Fd %d is not served, ingore stop", fd)
168 clnt = self.by_fd[fd]
169 log.info("Stop serving fd %d (IMEI %s)", clnt.sock.fileno(), clnt.imei)
171 if clnt.imei and self.by_imei[clnt.imei] == clnt: # could be replaced
172 del self.by_imei[clnt.imei]
177 ) -> Optional[List[Tuple[ProtoModule, Optional[str], float, Any, bytes]]]:
178 if fd not in self.by_fd:
179 log.debug("Client at fd %d gone, ingore event", fd)
181 clnt = self.by_fd[fd]
186 for when, peeraddr, packet in msgs:
187 assert clnt.pmod is not None
188 if clnt.imei is None:
189 imei = clnt.pmod.imei_from_packet(packet)
191 log.info("LOGIN from fd %d (IMEI %s)", fd, imei)
193 oldclnt = self.by_imei.get(clnt.imei)
194 if oldclnt is not None:
195 oldfd = oldclnt.sock.fileno()
196 log.info("Removing stale connection on fd %d", oldfd)
199 self.by_imei[clnt.imei] = clnt
200 result.append((clnt.pmod, clnt.imei, when, peeraddr, packet))
202 "Received from %s (IMEI %s): %s",
209 def response(self, resp: Resp) -> Optional[ProtoModule]:
210 if resp.imei in self.by_imei:
211 clnt = self.by_imei[resp.imei]
212 clnt.send(resp.packet)
215 log.info("Not connected (IMEI %s)", resp.imei)
219 def runserver(conf: ConfigParser, handle_hibernate: bool = True) -> None:
222 cast(ProtoModule, import_module("." + modnm, __package__))
223 for modnm in conf.get("collector", "protocols").split(",")
225 # Is this https://github.com/zeromq/pyzmq/issues/1627 still not fixed?!
226 zctx = zmq.Context() # type: ignore
227 zpub = zctx.socket(zmq.PUB) # type: ignore
228 zpull = zctx.socket(zmq.PULL) # type: ignore
229 oldmask = umask(0o117)
230 zpub.bind(conf.get("collector", "publishurl"))
231 zpull.bind(conf.get("collector", "listenurl"))
233 tcpl = socket(AF_INET6, SOCK_STREAM)
234 tcpl.setsockopt(SOL_SOCKET, SO_REUSEADDR, 1)
235 tcpl.bind(("", conf.getint("collector", "port")))
237 tcpfd = tcpl.fileno()
238 poller = zmq.Poller() # type: ignore
239 poller.register(zpull, flags=zmq.POLLIN)
240 poller.register(tcpfd, flags=zmq.POLLIN)
242 pollingfds: Set[int] = set()
245 tosend: List[Resp] = []
246 toadd: List[Tuple[socket, Any]] = []
247 events = poller.poll(1000)
248 for sk, fl in events:
252 msg = zpull.recv(zmq.NOBLOCK)
258 clntsock, clntaddr = tcpl.accept()
259 clntsock.setsockopt(SOL_SOCKET, SO_KEEPALIVE, 1)
260 toadd.append((clntsock, clntaddr))
261 elif fl & zmq.POLLIN:
262 received = clients.recv(sk)
264 log.debug("Terminal gone from fd %d", sk)
267 for pmod, imei, when, peeraddr, packet in received:
268 proto = pmod.proto_of_message(packet)
279 pmod.is_goodbye_packet(packet)
283 "Goodbye from fd %d (IMEI %s)",
288 respmsg = pmod.inline_response(packet)
289 if respmsg is not None:
291 Resp(imei=imei, when=when, packet=respmsg)
294 log.debug("Stray event: %s on socket %s", fl, sk)
295 # poll queue consumed, make changes now
297 log.debug("Sending to the client: %s", zmsg)
298 rpmod = clients.response(zmsg)
299 if rpmod is not None:
303 proto=rpmod.proto_of_message(zmsg.packet),
309 for fd in pollingfds - clients.fds():
310 poller.unregister(fd) # type: ignore
311 for clntsock, clntaddr in toadd:
312 fd = clients.add(clntsock, clntaddr)
313 for fd in clients.fds() - pollingfds:
314 poller.register(fd, flags=zmq.POLLIN)
315 pollingfds = clients.fds()
316 except KeyboardInterrupt:
319 zctx.destroy() # type: ignore
323 if __name__.endswith("__main__"):
324 runserver(common.init(log))