1 """ TCP server that communicates with terminals """
3 from logging import getLogger
5 from socket import socket, AF_INET6, SOCK_STREAM, SOL_SOCKET, SO_REUSEADDR
7 from struct import pack
11 from .gps303proto import (
18 from .zmsg import Bcast, Resp
20 log = getLogger("gps303/collector")
24 """Connected socket to the terminal plus buffer and metadata"""
26 def __init__(self, sock, addr):
33 log.debug("Closing fd %d (IMEI %s)", self.sock.fileno(), self.imei)
38 """Read from the socket and parse complete messages"""
40 segment = self.sock.recv(4096)
43 "Reading from fd %d (IMEI %s): %s",
49 if not segment: # Terminal has closed connection
51 "EOF reading from fd %d (IMEI %s)",
57 self.buffer += segment
60 framestart = self.buffer.find(b"xx")
61 if framestart == -1: # No frames, return whatever we have
63 if framestart > 0: # Should not happen, report
65 'Undecodable data "%s" from fd %d (IMEI %s)',
66 self.buffer[:framestart].hex(),
70 self.buffer = self.buffer[framestart:]
71 # At this point, buffer starts with a packet
72 if len(self.buffer) < 6: # no len and proto - cannot proceed
74 exp_end = self.buffer[2] + 3 # Expect '\r\n' here
76 # Length field can legitimeely be much less than the
77 # length of the packet (e.g. WiFi positioning), but
78 # it _should not_ be greater. Still sometimes it is.
79 # Luckily, not by too much: by maybe two or three bytes?
80 # Do this embarrassing hack to avoid accidental match
81 # of some binary data in the packet against '\r\n'.
83 frameend = self.buffer.find(b"\r\n", frameend)
84 if frameend >= (exp_end - 3): # Found realistic match
86 if frameend == -1: # Incomplete frame, return what we have
88 packet = self.buffer[2:frameend]
89 self.buffer = self.buffer[frameend + 2 :]
90 if proto_of_message(packet) == LOGIN.PROTO:
91 self.imei = parse_message(packet).imei
93 "LOGIN from fd %d (IMEI %s)", self.sock.fileno(), self.imei
95 msgs.append((when, self.addr, packet))
98 def send(self, buffer):
100 self.sock.send(b"xx" + buffer + b"\r\n")
103 "Sending to fd %d (IMEI %s): %s",
115 def add(self, clntsock, clntaddr):
116 fd = clntsock.fileno()
117 log.info("Start serving fd %d from %s", fd, clntaddr)
118 self.by_fd[fd] = Client(clntsock, clntaddr)
122 clnt = self.by_fd[fd]
123 log.info("Stop serving fd %d (IMEI %s)", clnt.sock.fileno(), clnt.imei)
126 del self.by_imei[clnt.imei]
130 clnt = self.by_fd[fd]
135 for when, peeraddr, packet in msgs:
136 if proto_of_message(packet) == LOGIN.PROTO: # Could do blindly...
137 self.by_imei[clnt.imei] = clnt
138 result.append((clnt.imei, when, peeraddr, packet))
140 "Received from %s (IMEI %s): %s",
147 def response(self, resp):
148 if resp.imei in self.by_imei:
149 self.by_imei[resp.imei].send(resp.packet)
151 log.info("Not connected (IMEI %s)", resp.imei)
156 zpub = zctx.socket(zmq.PUB)
157 zpull = zctx.socket(zmq.PULL)
158 oldmask = umask(0o117)
159 zpub.bind(conf.get("collector", "publishurl"))
160 zpull.bind(conf.get("collector", "listenurl"))
162 tcpl = socket(AF_INET6, SOCK_STREAM)
163 tcpl.setsockopt(SOL_SOCKET, SO_REUSEADDR, 1)
164 tcpl.bind(("", conf.getint("collector", "port")))
166 tcpfd = tcpl.fileno()
167 poller = zmq.Poller()
168 poller.register(zpull, flags=zmq.POLLIN)
169 poller.register(tcpfd, flags=zmq.POLLIN)
176 events = poller.poll(1000)
177 for sk, fl in events:
181 msg = zpull.recv(zmq.NOBLOCK)
187 clntsock, clntaddr = tcpl.accept()
188 topoll.append((clntsock, clntaddr))
189 elif fl & zmq.POLLIN:
190 received = clients.recv(sk)
192 log.debug("Terminal gone from fd %d", sk)
195 for imei, when, peeraddr, packet in received:
196 proto = proto_of_message(packet)
206 if proto == HIBERNATION.PROTO:
208 "HIBERNATION from fd %d (IMEI %s)",
213 respmsg = inline_response(packet)
214 if respmsg is not None:
216 Resp(imei=imei, when=when, packet=respmsg)
219 log.debug("Stray event: %s on socket %s", fl, sk)
220 # poll queue consumed, make changes now
222 poller.unregister(fd)
228 proto=proto_of_message(zmsg.packet),
234 log.debug("Sending to the client: %s", zmsg)
235 clients.response(zmsg)
236 for clntsock, clntaddr in topoll:
237 fd = clients.add(clntsock, clntaddr)
238 poller.register(fd, flags=zmq.POLLIN)
239 except KeyboardInterrupt:
243 if __name__.endswith("__main__"):
244 runserver(common.init(log))