1 """ TCP server that communicates with terminals """
3 from configparser import ConfigParser
4 from logging import getLogger
14 from struct import pack
16 from typing import Dict, List, Optional, Tuple
20 from .gps303proto import (
27 from .zmsg import Bcast, Resp
29 log = getLogger("gps303/collector")
35 """Connected socket to the terminal plus buffer and metadata"""
37 def __init__(self, sock: socket, addr: Tuple[str, int]) -> None:
41 self.imei: Optional[str] = None
43 def close(self) -> None:
44 log.debug("Closing fd %d (IMEI %s)", self.sock.fileno(), self.imei)
48 def recv(self) -> Optional[List[Tuple[float, Tuple[str, int], bytes]]]:
49 """Read from the socket and parse complete messages"""
51 segment = self.sock.recv(MAXBUFFER)
54 "Reading from fd %d (IMEI %s): %s",
60 if not segment: # Terminal has closed connection
62 "EOF reading from fd %d (IMEI %s)",
68 self.buffer += segment
69 if len(self.buffer) > MAXBUFFER:
70 # We are receiving junk. Let's drop it or we run out of memory.
71 log.warning("More than %d unparseable data, dropping", MAXBUFFER)
75 framestart = self.buffer.find(b"xx")
76 if framestart == -1: # No frames, return whatever we have
78 if framestart > 0: # Should not happen, report
80 'Undecodable data (%d) "%s" from fd %d (IMEI %s)',
82 self.buffer[:framestart][:64].hex(),
86 self.buffer = self.buffer[framestart:]
87 # At this point, buffer starts with a packet
88 if len(self.buffer) < 6: # no len and proto - cannot proceed
90 exp_end = self.buffer[2] + 3 # Expect '\r\n' here
92 # Length field can legitimeely be much less than the
93 # length of the packet (e.g. WiFi positioning), but
94 # it _should not_ be greater. Still sometimes it is.
95 # Luckily, not by too much: by maybe two or three bytes?
96 # Do this embarrassing hack to avoid accidental match
97 # of some binary data in the packet against '\r\n'.
99 frameend = self.buffer.find(b"\r\n", frameend + 1)
100 if frameend == -1 or frameend >= (
102 ): # Found realistic match or none
104 if frameend == -1: # Incomplete frame, return what we have
106 packet = self.buffer[2:frameend]
107 self.buffer = self.buffer[frameend + 2 :]
108 if len(packet) < 2: # frameend comes too early
109 log.warning("Packet too short: %s", packet)
111 msgs.append((when, self.addr, packet))
114 def send(self, buffer: bytes) -> None:
116 self.sock.send(b"xx" + buffer + b"\r\n")
119 "Sending to fd %d (IMEI %s): %s",
127 def __init__(self) -> None:
128 self.by_fd: Dict[int, Client] = {}
129 self.by_imei: Dict[str, Client] = {}
131 def add(self, clntsock: socket, clntaddr: Tuple[str, int]) -> int:
132 fd = clntsock.fileno()
133 log.info("Start serving fd %d from %s", fd, clntaddr)
134 self.by_fd[fd] = Client(clntsock, clntaddr)
137 def stop(self, fd: int) -> None:
138 clnt = self.by_fd[fd]
139 log.info("Stop serving fd %d (IMEI %s)", clnt.sock.fileno(), clnt.imei)
142 del self.by_imei[clnt.imei]
147 ) -> Optional[List[Tuple[Optional[str], float, Tuple[str, int], bytes]]]:
148 clnt = self.by_fd[fd]
153 for when, peeraddr, packet in msgs:
154 if proto_of_message(packet) == LOGIN.PROTO:
155 msg = parse_message(packet)
156 if isinstance(msg, LOGIN): # Can be unparseable
157 if clnt.imei is None:
160 "LOGIN from fd %d (IMEI %s)",
164 oldclnt = self.by_imei.get(clnt.imei)
165 if oldclnt is not None:
167 "Orphaning fd %d with the same IMEI",
168 oldclnt.sock.fileno(),
171 self.by_imei[clnt.imei] = clnt
174 "Login message from %s: %s, but client imei unfilled",
178 result.append((clnt.imei, when, peeraddr, packet))
180 "Received from %s (IMEI %s): %s",
187 def response(self, resp: Resp) -> None:
188 if resp.imei in self.by_imei:
189 self.by_imei[resp.imei].send(resp.packet)
191 log.info("Not connected (IMEI %s)", resp.imei)
194 def runserver(conf: ConfigParser, handle_hibernate: bool = True) -> None:
195 # Is this https://github.com/zeromq/pyzmq/issues/1627 still not fixed?!
196 zctx = zmq.Context() # type: ignore
197 zpub = zctx.socket(zmq.PUB) # type: ignore
198 zpull = zctx.socket(zmq.PULL) # type: ignore
199 oldmask = umask(0o117)
200 zpub.bind(conf.get("collector", "publishurl"))
201 zpull.bind(conf.get("collector", "listenurl"))
203 tcpl = socket(AF_INET6, SOCK_STREAM)
204 tcpl.setsockopt(SOL_SOCKET, SO_REUSEADDR, 1)
205 tcpl.bind(("", conf.getint("collector", "port")))
207 tcpfd = tcpl.fileno()
208 poller = zmq.Poller() # type: ignore
209 poller.register(zpull, flags=zmq.POLLIN)
210 poller.register(tcpfd, flags=zmq.POLLIN)
217 events = poller.poll(1000)
218 for sk, fl in events:
222 msg = zpull.recv(zmq.NOBLOCK)
228 clntsock, clntaddr = tcpl.accept()
229 clntsock.setsockopt(SOL_SOCKET, SO_KEEPALIVE, 1)
230 topoll.append((clntsock, clntaddr))
231 elif fl & zmq.POLLIN:
232 received = clients.recv(sk)
234 log.debug("Terminal gone from fd %d", sk)
237 for imei, when, peeraddr, packet in received:
238 proto = proto_of_message(packet)
248 if proto == HIBERNATION.PROTO and handle_hibernate:
250 "HIBERNATION from fd %d (IMEI %s)",
255 respmsg = inline_response(packet)
256 if respmsg is not None:
258 Resp(imei=imei, when=when, packet=respmsg)
261 log.debug("Stray event: %s on socket %s", fl, sk)
262 # poll queue consumed, make changes now
267 proto=proto_of_message(zmsg.packet),
273 log.debug("Sending to the client: %s", zmsg)
274 clients.response(zmsg)
276 poller.unregister(fd) # type: ignore
278 for clntsock, clntaddr in topoll:
279 fd = clients.add(clntsock, clntaddr)
280 poller.register(fd, flags=zmq.POLLIN)
281 except KeyboardInterrupt:
284 zctx.destroy() # type: ignore
288 if __name__.endswith("__main__"):
289 runserver(common.init(log))