]> www.average.org Git - loctrkd.git/blob - gps303/collector.py
make collector.py work
[loctrkd.git] / gps303 / collector.py
1 """ TCP server that communicates with terminals """
2
3 from getopt import getopt
4 from logging import getLogger, StreamHandler, DEBUG, INFO
5 from logging.handlers import SysLogHandler
6 from socket import socket, AF_INET, SOCK_STREAM, SOL_SOCKET, SO_REUSEADDR
7 from time import time
8 import sys
9 import zmq
10
11 from .config import readconfig
12 from .gps303proto import parse_message, HIBERNATION, LOGIN, set_config
13
14 CONF = "/etc/gps303.conf"
15
16 log = getLogger("gps303/collector")
17
18
19 class Bcast:
20     """Zmq message to broadcast what was received from the terminal"""
21     def __init__(self, imei, msg):
22         self.as_bytes = imei.encode() + msg.to_packet()
23
24
25 class Resp:
26     """Zmq message received from a third party to send to the terminal"""
27     def __init__(self, msg):
28         self.imei = msg[:16].decode()
29         self.payload = msg[16:]
30
31
32 class Client:
33     """Connected socket to the terminal plus buffer and metadata"""
34     def __init__(self, sock, addr):
35         self.sock = sock
36         self.addr = addr
37         self.buffer = b""
38         self.imei = None
39
40     def close(self):
41         log.debug("Closing fd %d (IMEI %s)", self.sock.fileno(), self.imei)
42         self.sock.close()
43         self.buffer = b""
44         self.imei = None
45
46     def recv(self):
47         """ Read from the socket and parse complete messages """
48         try:
49             segment = self.sock.recv(4096)
50         except OSError:
51             log.warning("Reading from fd %d (IMEI %s): %s",
52                     self.sock.fileno(), self.imei, e)
53             return None
54         if not segment:  # Terminal has closed connection
55             log.info("EOF reading from fd %d (IMEI %s)",
56                     self.sock.fileno(), self.imei)
57             return None
58         when = time()
59         self.buffer += segment
60         msgs = []
61         while True:
62             framestart = self.buffer.find(b"xx")
63             if framestart == -1:  # No frames, return whatever we have
64                 break
65             if framestart > 0:  # Should not happen, report
66                 log.warning("Undecodable data \"%s\" from fd %d (IMEI %s)",
67                         self.buffer[:framestart].hex(), self.sock.fileno(), self.imei)
68                 self.buffer = self.buffer[framestart:]
69             # At this point, buffer starts with a packet
70             frameend = self.buffer.find(b"\r\n", 4)
71             if frameend == -1:  # Incomplete frame, return what we have
72                 break
73             msg = parse_message(self.buffer[2:frameend])
74             self.buffer = self.buffer[frameend+2:]
75             if isinstance(msg, LOGIN):
76                 self.imei = msg.imei
77                 log.info("LOGIN from fd %d: IMEI %s",
78                         self.sock.fileno(), self.imei)
79             msgs.append(msg)
80         return msgs
81
82     def send(self, buffer):
83         try:
84             self.sock.send(b"xx" + buffer + b"\r\n")
85         except OSError as e:
86             log.error("Sending to fd %d (IMEI %s): %s",
87                     self.sock.fileno, self.imei, e)
88
89 class Clients:
90     def __init__(self):
91         self.by_fd = {}
92         self.by_imei = {}
93
94     def add(self, clntsock, clntaddr):
95         fd = clntsock.fileno()
96         self.by_fd[fd] = Client(clntsock, clntaddr)
97         return fd
98
99     def stop(self, fd):
100         clnt = self.by_fd[fd]
101         log.info("Stop serving fd %d (IMEI %s)", clnt.sock.fileno(), clnt.imei)
102         clnt.close()
103         if clnt.imei:
104             del self.by_imei[clnt.imei]
105         del self.by_fd[fd]
106
107     def recv(self, fd):
108         clnt = self.by_fd[fd]
109         msgs = clnt.recv()
110         result = []
111         for msg in msgs:
112             if isinstance(msg, LOGIN):
113                 self.by_imei[clnt.imei] = clnt
114             result.append((clnt.imei, msg))
115         return result
116
117     def response(self, resp):
118         if resp.imei in self.by_imei:
119             self.by_imei[resp.imei].send(resp.payload)
120
121
122 def runserver(opts, conf):
123     zctx = zmq.Context()
124     zpub = zctx.socket(zmq.PUB)
125     zpub.bind(conf.get("collector", "publishurl"))
126     zsub = zctx.socket(zmq.SUB)
127     zsub.connect(conf.get("collector", "listenurl"))
128     tcpl = socket(AF_INET, SOCK_STREAM)
129     tcpl.setsockopt(SOL_SOCKET, SO_REUSEADDR, 1)
130     tcpl.bind(("", conf.getint("collector", "port")))
131     tcpl.listen(5)
132     tcpfd = tcpl.fileno()
133     poller = zmq.Poller()
134     poller.register(zsub, flags=zmq.POLLIN)
135     poller.register(tcpfd, flags=zmq.POLLIN)
136     clients = Clients()
137     try:
138         while True:
139             tosend = []
140             topoll = []
141             tostop = []
142             events = poller.poll(10)
143             for sk, fl in events:
144                 if sk is zsub:
145                     while True:
146                         try:
147                             msg = zsub.recv(zmq.NOBLOCK)
148                             tosend.append(Resp(msg))
149                         except zmq.Again:
150                             break
151                 elif sk == tcpfd:
152                     clntsock, clntaddr = tcpl.accept()
153                     topoll.append((clntsock, clntaddr))
154                 else:
155                     for imei, msg in clients.recv(sk):
156                         zpub.send(Bcast(imei, msg).as_bytes)
157                         if msg is None or isinstance(msg, HIBERNATION):
158                             log.debug("HIBERNATION from fd %d", sk)
159                             tostop.append(sk)
160             # poll queue consumed, make changes now
161             for fd in tostop:
162                 poller.unregister(fd)
163                 clients.stop(fd)
164             for zmsg in tosend:
165                 clients.response(zmsg)
166             for clntsock, clntaddr in topoll:
167                 fd = clients.add(clntsock, clntaddr)
168                 poller.register(fd)
169     except KeyboardInterrupt:
170         pass
171
172
173 if __name__.endswith("__main__"):
174     opts, _ = getopt(sys.argv[1:], "c:d")
175     opts = dict(opts)
176     conf = readconfig(opts["-c"] if "-c" in opts else CONF)
177     if sys.stdout.isatty():
178         log.addHandler(StreamHandler(sys.stderr))
179     else:
180         log.addHandler(SysLogHandler(address="/dev/log"))
181     log.setLevel(DEBUG if "-d" in opts else INFO)
182     log.info("starting with options: %s", opts)
183     runserver(opts, conf)