1 """ TCP server that communicates with terminals """
3 from configparser import ConfigParser
4 from logging import getLogger
6 from socket import socket, AF_INET6, SOCK_STREAM, SOL_SOCKET, SO_REUSEADDR
7 from struct import pack
9 from typing import Dict, List, Optional, Tuple
13 from .gps303proto import (
20 from .zmsg import Bcast, Resp
22 log = getLogger("gps303/collector")
26 """Connected socket to the terminal plus buffer and metadata"""
28 def __init__(self, sock: socket, addr: Tuple[str, int]) -> None:
32 self.imei: Optional[str] = None
34 def close(self) -> None:
35 log.debug("Closing fd %d (IMEI %s)", self.sock.fileno(), self.imei)
39 def recv(self) -> Optional[List[Tuple[float, Tuple[str, int], bytes]]]:
40 """Read from the socket and parse complete messages"""
42 segment = self.sock.recv(4096)
45 "Reading from fd %d (IMEI %s): %s",
51 if not segment: # Terminal has closed connection
53 "EOF reading from fd %d (IMEI %s)",
59 self.buffer += segment
62 framestart = self.buffer.find(b"xx")
63 if framestart == -1: # No frames, return whatever we have
65 if framestart > 0: # Should not happen, report
67 'Undecodable data "%s" from fd %d (IMEI %s)',
68 self.buffer[:framestart].hex(),
72 self.buffer = self.buffer[framestart:]
73 # At this point, buffer starts with a packet
74 if len(self.buffer) < 6: # no len and proto - cannot proceed
76 exp_end = self.buffer[2] + 3 # Expect '\r\n' here
78 # Length field can legitimeely be much less than the
79 # length of the packet (e.g. WiFi positioning), but
80 # it _should not_ be greater. Still sometimes it is.
81 # Luckily, not by too much: by maybe two or three bytes?
82 # Do this embarrassing hack to avoid accidental match
83 # of some binary data in the packet against '\r\n'.
85 frameend = self.buffer.find(b"\r\n", frameend)
86 if frameend >= (exp_end - 3): # Found realistic match
88 if frameend == -1: # Incomplete frame, return what we have
90 packet = self.buffer[2:frameend]
91 self.buffer = self.buffer[frameend + 2 :]
92 if proto_of_message(packet) == LOGIN.PROTO:
93 self.imei = parse_message(packet).imei
95 "LOGIN from fd %d (IMEI %s)", self.sock.fileno(), self.imei
97 msgs.append((when, self.addr, packet))
100 def send(self, buffer: bytes) -> None:
102 self.sock.send(b"xx" + buffer + b"\r\n")
105 "Sending to fd %d (IMEI %s): %s",
113 def __init__(self) -> None:
114 self.by_fd: Dict[int, Client] = {}
115 self.by_imei: Dict[str, Client] = {}
117 def add(self, clntsock: socket, clntaddr: Tuple[str, int]) -> int:
118 fd = clntsock.fileno()
119 log.info("Start serving fd %d from %s", fd, clntaddr)
120 self.by_fd[fd] = Client(clntsock, clntaddr)
123 def stop(self, fd: int) -> None:
124 clnt = self.by_fd[fd]
125 log.info("Stop serving fd %d (IMEI %s)", clnt.sock.fileno(), clnt.imei)
128 del self.by_imei[clnt.imei]
131 def recv(self, fd: int) -> Optional[List[Tuple[Optional[str], float, Tuple[str, int], bytes]]]:
132 clnt = self.by_fd[fd]
137 for when, peeraddr, packet in msgs:
138 if proto_of_message(packet) == LOGIN.PROTO: # Could do blindly...
140 self.by_imei[clnt.imei] = clnt
142 log.warning("Login message from %s: %s, but client imei unfilled", peeraddr, packet)
143 result.append((clnt.imei, when, peeraddr, packet))
145 "Received from %s (IMEI %s): %s",
152 def response(self, resp: Resp) -> None:
153 if resp.imei in self.by_imei:
154 self.by_imei[resp.imei].send(resp.packet)
156 log.info("Not connected (IMEI %s)", resp.imei)
159 def runserver(conf: ConfigParser) -> None:
160 # Is this https://github.com/zeromq/pyzmq/issues/1627 still not fixed?!
161 zctx = zmq.Context() # type: ignore
162 zpub = zctx.socket(zmq.PUB) # type: ignore
163 zpull = zctx.socket(zmq.PULL) # type: ignore
164 oldmask = umask(0o117)
165 zpub.bind(conf.get("collector", "publishurl"))
166 zpull.bind(conf.get("collector", "listenurl"))
168 tcpl = socket(AF_INET6, SOCK_STREAM)
169 tcpl.setsockopt(SOL_SOCKET, SO_REUSEADDR, 1)
170 tcpl.bind(("", conf.getint("collector", "port")))
172 tcpfd = tcpl.fileno()
173 poller = zmq.Poller() # type: ignore
174 poller.register(zpull, flags=zmq.POLLIN)
175 poller.register(tcpfd, flags=zmq.POLLIN)
182 events = poller.poll(1000)
183 for sk, fl in events:
187 msg = zpull.recv(zmq.NOBLOCK)
193 clntsock, clntaddr = tcpl.accept()
194 topoll.append((clntsock, clntaddr))
195 elif fl & zmq.POLLIN:
196 received = clients.recv(sk)
198 log.debug("Terminal gone from fd %d", sk)
201 for imei, when, peeraddr, packet in received:
202 proto = proto_of_message(packet)
212 if proto == HIBERNATION.PROTO:
214 "HIBERNATION from fd %d (IMEI %s)",
219 respmsg = inline_response(packet)
220 if respmsg is not None:
222 Resp(imei=imei, when=when, packet=respmsg)
225 log.debug("Stray event: %s on socket %s", fl, sk)
226 # poll queue consumed, make changes now
228 poller.unregister(fd) # type: ignore
234 proto=proto_of_message(zmsg.packet),
240 log.debug("Sending to the client: %s", zmsg)
241 clients.response(zmsg)
242 for clntsock, clntaddr in topoll:
243 fd = clients.add(clntsock, clntaddr)
244 poller.register(fd, flags=zmq.POLLIN)
245 except KeyboardInterrupt:
249 if __name__.endswith("__main__"):
250 runserver(common.init(log))