-
Notifications
You must be signed in to change notification settings - Fork 1
/
client.py
247 lines (202 loc) · 8.92 KB
/
client.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
#!/usr/bin/env python
# -*- coding: utf-8 -*-
from __future__ import print_function, unicode_literals
import argparse
import logging
import coloredlogs
import yaml
import os
import sys
import geofence
import storage
import datex2
import time
from interchange import NordicWayIC, ConnectionError
from util import slack_notify
log = logging.getLogger("geofencebroker")
def init_logging(debug_env_var):
field_style_override = coloredlogs.DEFAULT_FIELD_STYLES
level_style_override = coloredlogs.DEFAULT_LEVEL_STYLES
logging_level = 'INFO'
log.setLevel(logging.INFO)
if os.environ.get(debug_env_var):
logging_level = 'DEBUG'
log.setLevel(logging.DEBUG)
field_style_override['levelname'] = {"color": "magenta", "bold": True}
level_style_override['debug'] = {"color": "blue"}
coloredlogs.install(level=logging_level,
fmt="%(asctime)s %(name)s %(levelname)s %(message)s",
level_styles=level_style_override,
field_styles=field_style_override)
qpid_log = logging.getLogger("qpid.messaging")
qpid_log.setLevel(logging.INFO)
# ch = logging.StreamHandler(stream=sys.stdout)
# formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
# ch.setFormatter(formatter)
# log.addHandler(ch)
if __name__ == '__main__':
parser = argparse.ArgumentParser()
parser.add_argument("-conf", "--config", help="Config file that specifies all input parameters", default=None)
parser.add_argument("-b", "--broker_url", help="AMQP schemed url to the broker/AMQP server.", default=None)
parser.add_argument("-v", "--verbose", help="Enable verbose logging / printing", default=False)
parser.add_argument("-s", "--sender", help="Name of the sender queue", default=None)
parser.add_argument("-r", "--receiver", help="Name of the receiver queue", default=None)
parser.add_argument("-k", "--ssl-keyfile", help="SSL key file")
parser.add_argument("-c", "--ssl-certfile", help="SSL cert file")
parser.add_argument("-u", "--username", help="Username", default=None)
parser.add_argument("-p", "--password", help="Password", default=None)
parser.add_argument("-t", "--timeout", type=int,
help="Timeout in seconds before checking NVDB for geofence updates", default=None)
args = parser.parse_args()
cfg = {}
init_logging('DEBUG')
if args.config:
log.info("Using config file")
if not os.path.exists(args.config):
log.error("Config file was not found..")
sys.exit(1)
try:
with open(args.config, 'r') as f:
cfg = yaml.load(f)
except yaml.scanner.ScannerError as se:
log.error("Error parsing config file: {}".format(args.config))
sys.exit(1)
required_parms = []
cfg_get = lambda x: cfg.get(x, False)
map(required_parms.append, [cfg_get("broker_url"), cfg_get("sender"), cfg_get("receiver")])
if not all(required_parms):
log.error("Missing required parameters from config file!")
sys.exit(1)
else:
required_parms = []
map(required_parms.append,
[args.broker_url, args.sender, args.receiver])
log.debug("required_parms: {}".format(required_parms))
if not all(required_parms):
log.error("Missing required parameters!")
sys.exit(1)
if args.verbose:
cfg.update({"verbose": True})
if args.timeout:
cfg.update({"timeout": args.timeout})
options = {"ssl_skip_hostname_check": True}
if cfg.get("ssl_keyfile", False):
options.update({"ssl_keyfile": cfg.get("ssl_keyfile")})
if cfg.get("ssl_certfile", False):
options.update({"ssl_certfile": cfg.get("ssl_certfile")})
if cfg.get("broker_url", "").startswith("amqps"):
# Encrypted amqps session
if not all([cfg.get("ssl_keyfile"), cfg.get("ssl_certfile")]):
log.error("""Broker URL uses TLS/SSL.
Therefore you need to specify SSL cert and key.""")
sys.exit(1)
# Enable reconnection option
options.update({
"reconnect": True
})
log.info("Connecting to {broker_url}".format(**cfg))
log.info(" sender: {sender}, receiver: {receiver}".format(**cfg))
ic = NordicWayIC(cfg.get("broker_url"),
cfg.get("sender"),
cfg.get("receiver"),
cfg.get("username"),
cfg.get("password"),
options)
log.debug(ic)
# try:
# ic.connect()
# except ConnectionError as e:
# log.error("Unable to connect to {}".format(cfg.get("broker_url")))
# log.error(e)
# sys.exit(1)
# if not ic.connection.opened():
# log.error("Unable to connect!")
# sys.exit(1)
sleep_time = cfg.get("timeout")
log.debug("Sleeping for {} seconds between each check.".format(sleep_time))
# hack to create centroids if missing - a one-time operation!
storage.fix_centroid()
slack_url = cfg.get("slack_webhook_url", None)
slack_notify(
"Geofence is up and running! Will check periodically every {} second".format(sleep_time),
slack_url)
# Main loop
while True:
fences = geofence.fetch_objects()
if not fences or fences["metadata"].get("returnert", 0) == 0:
time.sleep(sleep_time)
continue
# TODO: Check if returned JSON has paging. If so, fetch the rest of
# the geofence objects
vegobjekt_ids = []
try:
ic.connect()
log.debug("Connect to interchange.")
for fence in fences.get("objekter"):
# Sample all vegobjekt IDs to check if there are anyone that
# has been deleted from NVDB
vegobjekt_ids.append(int(fence.get("id", 0)))
if not storage.exists(fence):
datex_obj = datex2.create_doc(fence)
#datex_obj.name = unicode("TestÆØÅ-New")
msg = u"New geofence: id={}, version={}, name={}".format(
fence.get("id"), datex_obj.version, datex_obj.name)
log.info(msg)
try:
slack_notify(msg, slack_url)
except Exception:
log.warn("Unable to send slack notification")
try:
ic.send_obj(datex_obj)
except ConnectionError as ce:
raise ce
else:
storage.add(fence)
else:
if storage.is_modified(fence):
datex_obj = datex2.create_doc(fence)
# datex_obj.name = unicode("TestÆØÅ")
msg = u"Modified geofence: message: id={}, version={}, name={}".format(
fence.get("id"), datex_obj.version, datex_obj.name)
log.info(msg)
slack_notify(msg, slack_url)
try:
ic.send_obj(datex_obj)
except ConnectionError as ce:
raise ce
else:
storage.update(fence)
# else:
# log.debug("geofence is already in db and has not been updated.")
# Find deleted vegobjekter in NVDB by getting a list of ID's from our
# cache database and list all ID's missing in our JSON from NVDB.
vegobjekter = storage.vegobjekter().all()
for v in vegobjekter:
# log.debug("inspecting v: {}".format(v.get("id")))
if v.get("id") not in vegobjekt_ids:
msg = "Vegobjekt with ID '{}' removed from NVDB: {}".format(v.get("id"), v)
log.warn(msg)
slack_notify(msg, slack_url)
datex_obj = datex2.create_delete_doc_from_db(v)
ic.send_obj(datex_obj)
log.debug(datex_obj)
storage.delete(v.get("id"))
log.warn("Delete geofence id: {}".format(v.get("id")))
except ConnectionError:
# Interchange lost its connection
log.debug("Interchange connection error. Trying to re-connect. URI: {}".format(cfg.get("broker_url")))
#while True:
# try:
# ic.connect()
# log.debug("Successfully re-connected to {}".format(cfg.get("broker_url")))
# except ConnectionError:
# log.debug("")
# time.sleep(5)
# continue
# break
else:
log.debug("Disconnect interchange.")
ic.close()
time.sleep(sleep_time)
ic.close()
log.info("Shutdown.. See ya!")