1 """ TCP server that communicates with terminals """
3 from logging import getLogger
4 from socket import socket, AF_INET6, SOCK_STREAM, SOL_SOCKET, SO_REUSEADDR
6 from struct import pack
10 from .gps303proto import (
17 from .zmsg import Bcast, Resp
19 log = getLogger("gps303/collector")
23 """Connected socket to the terminal plus buffer and metadata"""
25 def __init__(self, sock, addr):
32 log.debug("Closing fd %d (IMEI %s)", self.sock.fileno(), self.imei)
37 """Read from the socket and parse complete messages"""
39 segment = self.sock.recv(4096)
42 "Reading from fd %d (IMEI %s): %s",
48 if not segment: # Terminal has closed connection
50 "EOF reading from fd %d (IMEI %s)",
56 self.buffer += segment
59 framestart = self.buffer.find(b"xx")
60 if framestart == -1: # No frames, return whatever we have
62 if framestart > 0: # Should not happen, report
64 'Undecodable data "%s" from fd %d (IMEI %s)',
65 self.buffer[:framestart].hex(),
69 self.buffer = self.buffer[framestart:]
70 # At this point, buffer starts with a packet
71 frameend = self.buffer.find(b"\r\n", 4)
72 if frameend == -1: # Incomplete frame, return what we have
74 packet = self.buffer[2:frameend]
75 self.buffer = self.buffer[frameend + 2 :]
76 if proto_of_message(packet) == LOGIN.PROTO:
77 self.imei = parse_message(packet).imei
79 "LOGIN from fd %d (IMEI %s)", self.sock.fileno(), self.imei
81 msgs.append((when, self.addr, packet))
84 def send(self, buffer):
86 self.sock.send(b"xx" + buffer + b"\r\n")
89 "Sending to fd %d (IMEI %s): %s",
101 def add(self, clntsock, clntaddr):
102 fd = clntsock.fileno()
103 log.info("Start serving fd %d from %s", fd, clntaddr)
104 self.by_fd[fd] = Client(clntsock, clntaddr)
108 clnt = self.by_fd[fd]
109 log.info("Stop serving fd %d (IMEI %s)", clnt.sock.fileno(), clnt.imei)
112 del self.by_imei[clnt.imei]
116 clnt = self.by_fd[fd]
121 for when, peeraddr, packet in msgs:
122 if proto_of_message(packet) == LOGIN.PROTO: # Could do blindly...
123 self.by_imei[clnt.imei] = clnt
124 result.append((clnt.imei, when, peeraddr, packet))
126 "Received from %s (IMEI %s): %s",
133 def response(self, resp):
134 if resp.imei in self.by_imei:
135 self.by_imei[resp.imei].send(resp.packet)
137 log.info("Not connected (IMEI %s)", resp.imei)
142 zpub = zctx.socket(zmq.PUB)
143 zpub.bind(conf.get("collector", "publishurl"))
144 zpull = zctx.socket(zmq.PULL)
145 zpull.bind(conf.get("collector", "listenurl"))
146 tcpl = socket(AF_INET6, SOCK_STREAM)
147 tcpl.setsockopt(SOL_SOCKET, SO_REUSEADDR, 1)
148 tcpl.bind(("", conf.getint("collector", "port")))
150 tcpfd = tcpl.fileno()
151 poller = zmq.Poller()
152 poller.register(zpull, flags=zmq.POLLIN)
153 poller.register(tcpfd, flags=zmq.POLLIN)
160 events = poller.poll(1000)
161 for sk, fl in events:
165 msg = zpull.recv(zmq.NOBLOCK)
166 tosend.append(Resp(msg))
170 clntsock, clntaddr = tcpl.accept()
171 topoll.append((clntsock, clntaddr))
172 elif fl & zmq.POLLIN:
173 received = clients.recv(sk)
176 "Terminal gone from fd %d (IMEI %s)", sk, imei
180 for imei, when, peeraddr, packet in received:
181 proto = proto_of_message(packet)
191 if proto == HIBERNATION.PROTO:
193 "HIBERNATION from fd %d (IMEI %s)",
198 respmsg = inline_response(packet)
199 if respmsg is not None:
201 Resp(imei=imei, packet=respmsg)
204 log.debug("Stray event: %s on socket %s", fl, sk)
205 # poll queue consumed, make changes now
207 poller.unregister(fd)
210 log.debug("Sending to the client: %s", zmsg)
211 clients.response(zmsg)
212 for clntsock, clntaddr in topoll:
213 fd = clients.add(clntsock, clntaddr)
214 poller.register(fd, flags=zmq.POLLIN)
215 except KeyboardInterrupt:
219 if __name__.endswith("__main__"):
220 runserver(common.init(log))